mqx

package
v1.11.1 Latest Latest
Warning

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

Go to latest
Published: Jul 29, 2026 License: MIT Imports: 9 Imported by: 0

Documentation

Index

Constants

View Source
const (
	HeaderTenantID = "x-tenant-id"
	HeaderUserID   = "x-user-id"
	HeaderBizID    = "x-biz-id"
)

Variables

This section is empty.

Functions

func ConsumeMessage added in v1.11.1

func ConsumeMessage(ctx context.Context, consumer mq.Consumer) (context.Context, *mq.Message, error)

ConsumeMessage 消费原始 MQ 消息,并返回从消息头恢复后的业务上下文。

func ExtractContext added in v1.11.1

func ExtractContext(ctx context.Context, msg *mq.Message) context.Context

ExtractContext 从 mq.Message.Header 中提取租户 ID、用户 ID 和业务 ID,重新注入并返回新的 Context

func InjectContext added in v1.11.1

func InjectContext(ctx context.Context, msg *mq.Message)

InjectContext 提取 Context 中的租户 ID、用户 ID 和业务 ID,注入到 mq.Message.Header 中

func ProduceMessage added in v1.11.1

func ProduceMessage(ctx context.Context, producer mq.Producer, msg *mq.Message) (*mq.ProducerResult, error)

ProduceMessage 注入上下文元数据后发送原始 MQ 消息。

Types

type ConsumeFunc added in v1.11.1

type ConsumeFunc func(ctx context.Context, message *mq.Message) error

type Consumer added in v1.11.1

type Consumer struct {
	// contains filtered or unexported fields
}

func NewConsumer added in v1.11.1

func NewConsumer(name string, mq mq.MQ, topic string) *Consumer

func (*Consumer) Name added in v1.11.1

func (c *Consumer) Name() string

func (*Consumer) Start added in v1.11.1

func (c *Consumer) Start(ctx context.Context, consumeFunc ConsumeFunc) error

func (*Consumer) Stop added in v1.11.1

func (c *Consumer) Stop() error

type GeneralProducer

type GeneralProducer[T any] struct {
	// contains filtered or unexported fields
}

func NewGeneralProducer

func NewGeneralProducer[T any](q mq.MQ, topic string) (*GeneralProducer[T], error)

func (*GeneralProducer[T]) Produce

func (p *GeneralProducer[T]) Produce(ctx context.Context, evt T) error

type MultipleProducer

type MultipleProducer[T any] struct {
	// contains filtered or unexported fields
}

MultipleProducer 管理多个 GeneralProducer

func NewMultipleProducer

func NewMultipleProducer[T any](mq mq.MQ) *MultipleProducer[T]

NewMultipleProducer 创建一个新的 ProducerManager

func (*MultipleProducer[T]) AddProducer

func (pm *MultipleProducer[T]) AddProducer(topic string) error

AddProducer 添加一个新的 producer 监听指定的 topic

func (*MultipleProducer[T]) Produce

func (pm *MultipleProducer[T]) Produce(ctx context.Context, topic string, evt T) error

Produce 发送消息到指定的 topic

type Producer

type Producer[T any] interface {
	Produce(ctx context.Context, evt T) error
}

Jump to

Keyboard shortcuts

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