Documentation
¶
Index ¶
Constants ¶
View Source
const ( DefaultMaxQueueConcurrencyNum = 1 DefaultMaxQueueLen = 10000 )
View Source
const (
TaskQueueKey = "Platform:Queue:Task:%s"
)
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type ConsumerFunc ¶
type ConsumerFunc func(Delivery)
type OptionFunc ¶
type OptionFunc func(*TaskQueueEntity)
func ConfigConsumerFunc ¶
func ConfigConsumerFunc(f ConsumerFunc) OptionFunc
func ConfigMaxQueueConcurrencyNum ¶
func ConfigMaxQueueConcurrencyNum(num int64) OptionFunc
func ConfigMaxQueueLen ¶
func ConfigMaxQueueLen(num int64) OptionFunc
func ConfigName ¶
func ConfigName(name string) OptionFunc
func ConfigRedisConn ¶
func ConfigRedisConn(redisConn *redis.Client) OptionFunc
type Queue ¶
type Queue interface {
// Publish 推送消息
Publish(payload string) error
// AddConsumerFunc 添加消费者函数
AddConsumerFunc(consumerFunc ConsumerFunc)
// StartConsuming 开启消费
StartConsuming() error
}
func NewTaskQueueEntity ¶
func NewTaskQueueEntity(options ...OptionFunc) (Queue, error)
type TaskQueueEntity ¶
type TaskQueueEntity struct {
Name string // 名称
QueueKey string // 队列Key,根据名称计算
MaxQueueConcurrencyNum int64 // 并发处理数
RedisConn *redis.Client // Redis 连接
ConsumerFuncs []ConsumerFunc // 消息处理方法
MaxQueueLen int64 // 队列最大长度
// contains filtered or unexported fields
}
TaskQueueEntity 队列参数结构体
func (*TaskQueueEntity) AddConsumerFunc ¶
func (q *TaskQueueEntity) AddConsumerFunc(consumerFunc ConsumerFunc)
func (*TaskQueueEntity) Publish ¶
func (q *TaskQueueEntity) Publish(payload string) error
func (*TaskQueueEntity) StartConsuming ¶
func (q *TaskQueueEntity) StartConsuming() error
Click to show internal directories.
Click to hide internal directories.