Documentation
¶
Index ¶
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func Init ¶
func Init(cfg KafkaConfig) error
Init initializes the Kafka producer singleton with config.
Types ¶
type ConsumerConfig ¶
type ConsumerConfig struct {
Name string `toml:"name"` // 消费者名称(用于日志)
Group string `toml:"group"` // 消费组
Topics []string `toml:"topics"` // 订阅的 topics
}
ConsumerConfig 单个消费者配置
type InMemoryQueue ¶
type InMemoryQueue struct {
// contains filtered or unexported fields
}
InMemoryQueue 内存消息队列(用于测试和简单场景)
func (*InMemoryQueue) GetMessages ¶
func (q *InMemoryQueue) GetMessages(topic string) [][]byte
GetMessages 获取指定 topic 的所有消息(用于测试)
type KafkaConfig ¶
type KafkaConfig struct {
Brokers []string `toml:"brokers"`
Consumers []ConsumerConfig `toml:"consumers"`
Enabled bool `toml:"enabled"`
}
KafkaConfig Kafka 配置
type KafkaConsumer ¶
type KafkaConsumer struct {
// contains filtered or unexported fields
}
KafkaConsumer Kafka 消费者
func NewKafkaConsumer ¶
func NewKafkaConsumer(brokers []string, config ConsumerConfig, handler MessageHandler) (*KafkaConsumer, error)
NewKafkaConsumer 创建 Kafka 消费者
type KafkaProducer ¶
type KafkaProducer struct {
// contains filtered or unexported fields
}
KafkaProducer Kafka 生产者
func NewKafkaProducer ¶
func NewKafkaProducer(config KafkaConfig) (*KafkaProducer, error)
NewKafkaProducer 创建 Kafka 生产者
func NewQueue ¶
func NewQueue() *KafkaProducer
NewQueue returns the singleton Kafka producer instance. Returns nil if Kafka is not enabled or not initialized.
type MessageHandler ¶
MessageHandler 消息处理函数
Click to show internal directories.
Click to hide internal directories.