queue

package
v0.0.31 Latest Latest
Warning

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

Go to latest
Published: Sep 11, 2026 License: MIT Imports: 14 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

View Source
var ProviderSet = wire.NewSet(NewServer)

ProviderSet 创建队列消费服务,并注册 Core 与宿主提供的消费者。

Functions

func AddQueue

func AddQueue(queueName queue.Stream, data any) bool

AddQueue 向运行时队列追加异步消息。

func Decode

func Decode[T any](message data.Message) (*T, error)

Decode 从队列消息的 data 字段解码业务对象。

Types

type Consumer added in v0.0.6

type Consumer struct {
	// Stream 是消费者监听的队列流名称。
	Stream queueTransport.Stream
	// Handler 是收到队列消息后的处理函数。
	Handler queueData.ConsumerFunc
}

Consumer 描述宿主提供的一项队列消费者。

type Consumers added in v0.0.6

type Consumers []Consumer

Consumers 聚合宿主提供的队列消费者。

type Server

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

Server 将 Core 队列服务接入 Kratos 应用生命周期。

func NewServer

func NewServer(queue kitQueue.Queue, jobStore data.JobStore, logPipeline *biz.LogPipeline, consumers Consumers) (*Server, error)

NewServer 创建队列服务,注册 Core 和宿主提供的队列消费者。

func (*Server) Start

func (s *Server) Start(ctx context.Context) error

Start 启动队列消费服务。

func (*Server) Stop

func (s *Server) Stop(ctx context.Context) error

Stop 停止队列消费服务。

Jump to

Keyboard shortcuts

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