mq

package
v0.0.1 Latest Latest
Warning

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

Go to latest
Published: Dec 24, 2025 License: MIT Imports: 6 Imported by: 0

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 NewInMemoryQueue

func NewInMemoryQueue() *InMemoryQueue

NewInMemoryQueue 创建内存消息队列

func (*InMemoryQueue) Close

func (q *InMemoryQueue) Close() error

Close 关闭

func (*InMemoryQueue) GetMessages

func (q *InMemoryQueue) GetMessages(topic string) [][]byte

GetMessages 获取指定 topic 的所有消息(用于测试)

func (*InMemoryQueue) Publish

func (q *InMemoryQueue) Publish(topic string, message []byte) error

Publish 发布消息(同步处理)

func (*InMemoryQueue) Subscribe

func (q *InMemoryQueue) Subscribe(topic string, handler func([]byte) error) error

Subscribe 订阅 topic

type KafkaConfig

type KafkaConfig struct {
	Brokers   []string         `toml:"brokers"`
	Consumers []ConsumerConfig `toml:"consumers"`
	Enabled   bool             `toml:"enabled"`
}

KafkaConfig Kafka 配置

func (*KafkaConfig) Validate

func (c *KafkaConfig) Validate() error

Validate 验证配置

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 消费者

func (*KafkaConsumer) Start

func (c *KafkaConsumer) Start(ctx context.Context) error

Start 启动消费者

func (*KafkaConsumer) Stop

func (c *KafkaConsumer) Stop() error

Stop 停止消费者

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.

func (*KafkaProducer) Close

func (p *KafkaProducer) Close() error

Close 关闭生产者

func (*KafkaProducer) Publish

func (p *KafkaProducer) Publish(topic string, message []byte) error

Publish 发布消息

func (*KafkaProducer) Subscribe

func (p *KafkaProducer) Subscribe(topic string, handler func(message []byte) error) error

Subscribe 订阅(Producer 不支持,仅用于满足 MessageQueue 接口)

type MessageHandler

type MessageHandler func(ctx context.Context, topic string, message []byte) error

MessageHandler 消息处理函数

type MessageQueue

type MessageQueue interface {
	Publish(topic string, message []byte) error
	Subscribe(topic string, handler func(message []byte) error) error
	Close() error
}

MessageQueue 消息队列接口

Jump to

Keyboard shortcuts

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