Documentation
¶
Index ¶
- Constants
- func ConsumeMessage(ctx context.Context, consumer mq.Consumer) (context.Context, *mq.Message, error)
- func ExtractContext(ctx context.Context, msg *mq.Message) context.Context
- func InjectContext(ctx context.Context, msg *mq.Message)
- func ProduceMessage(ctx context.Context, producer mq.Producer, msg *mq.Message) (*mq.ProducerResult, error)
- type ConsumeFunc
- type Consumer
- type GeneralProducer
- type MultipleProducer
- type Producer
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
ExtractContext 从 mq.Message.Header 中提取租户 ID、用户 ID 和业务 ID,重新注入并返回新的 Context
func InjectContext ¶ added in v1.11.1
InjectContext 提取 Context 中的租户 ID、用户 ID 和业务 ID,注入到 mq.Message.Header 中
Types ¶
type ConsumeFunc ¶ added in v1.11.1
type Consumer ¶ added in v1.11.1
type Consumer struct {
// contains filtered or unexported fields
}
type GeneralProducer ¶
type GeneralProducer[T any] struct { // contains filtered or unexported fields }
func NewGeneralProducer ¶
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
Click to show internal directories.
Click to hide internal directories.