kafka

package
v0.0.3 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 18, 2025 License: MIT Imports: 6 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func Errorf

func Errorf(msg string, args ...interface{})

func Infof

func Infof(msg string, args ...interface{})

Types

type Consumer

type Consumer struct {
	*kafka.Reader
}

func DisconveryConsumer

func DisconveryConsumer(topic, groupID string, _opts ...OptionConsumer) (*Consumer, error)

func DisconveryConsumerAppExclusive

func DisconveryConsumerAppExclusive(app, topic, groupID string, _opts ...OptionConsumer) (*Consumer, error)

func (*Consumer) Ack

func (c *Consumer) Ack(ctx context.Context, msgs ...Message) error

func (*Consumer) Receive

func (c *Consumer) Receive(ctx context.Context) (Message, error)

Receive 读取一条消息, 并且自动 commit.

func (*Consumer) ReceiveNotCommit

func (c *Consumer) ReceiveNotCommit(ctx context.Context) (Message, error)

type Message

type Message = kafka.Message

type OptionConsumer

type OptionConsumer func(*kafka.ReaderConfig)

func OptionConsumerAddr

func OptionConsumerAddr(addrs []string) OptionConsumer

func OptionConsumerCommitInterval

func OptionConsumerCommitInterval(dur time.Duration) OptionConsumer

func OptionConsumerGroupID

func OptionConsumerGroupID(groupID string) OptionConsumer

func OptionConsumerQueueCapacity

func OptionConsumerQueueCapacity(cap int) OptionConsumer

func OptionConsumerSync

func OptionConsumerSync() OptionConsumer

func OptionConsumerTopic

func OptionConsumerTopic(topic string) OptionConsumer

type OptionProducer

type OptionProducer func(*kafka.Writer)

func OptionProducerAddr

func OptionProducerAddr(addrs []string) OptionProducer

func OptionProducerAsync

func OptionProducerAsync() OptionProducer

func OptionProducerBalancer

func OptionProducerBalancer(balance kafka.Balancer) OptionProducer

func OptionProducerBatchSize

func OptionProducerBatchSize(size int) OptionProducer

func OptionProducerReadTimeout

func OptionProducerReadTimeout(dur time.Duration) OptionProducer

func OptionProducerRequiredAcks

func OptionProducerRequiredAcks(level kafka.RequiredAcks) OptionProducer

func OptionProducerTopic

func OptionProducerTopic(topic string) OptionProducer

func OptionProducerWithErrorLog

func OptionProducerWithErrorLog(logger kafka.LoggerFunc) OptionProducer

func OptionProducerWriteTimeout

func OptionProducerWriteTimeout(dur time.Duration) OptionProducer

type Producer

type Producer struct {
	*kafka.Writer
}

func DisconveryProducer

func DisconveryProducer(topic string, _opts ...OptionProducer) (*Producer, error)

func DisconveryProducerAppExclusive

func DisconveryProducerAppExclusive(app, topic string, _opts ...OptionProducer) (*Producer, error)

func (*Producer) Close

func (p *Producer) Close(ctx context.Context) error

func (*Producer) Send

func (p *Producer) Send(ctx context.Context, msgs ...Message) error

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL