Documentation
¶
Index ¶
- Constants
- Variables
- func DemoCoA(ctx context.Context)
- func DemoCoB(ctx context.Context)
- func GetContextRunId(ctx context.Context) int
- func GetContextRunName(ctx context.Context) string
- func SetContextRunId(ctx context.Context, index int) context.Context
- func SetContextRunName(ctx context.Context, index string) context.Context
- type CoFunc
- type ConProcessModelActuator
- type CtlRater
- type EventBus
- type FailAllFailCoMag
- type FragmentProcessor
- type ProcessMultiSender
- type ProcessSender
- type PubSubActuator
- func (bus *PubSubActuator[T]) Publish(topic string, msgs ...T) error
- func (bus *PubSubActuator[T]) PublishByFn(topic string, msgs ...T) error
- func (bus *PubSubActuator[T]) Start() error
- func (bus *PubSubActuator[T]) Stop() error
- func (bus *PubSubActuator[T]) StopAsync() error
- func (bus *PubSubActuator[T]) Subscribe(topic string) (<-chan T, error)
- func (bus *PubSubActuator[T]) SubscribeByFn(topic string, fn func(int, T)) error
- func (bus *PubSubActuator[T]) WaitStop()
- type PubSubConfig
- type PubSubOption
- type StopNotifyWorker
- type StopWorker
Constants ¶
const ContextRunId = "__ID__"
const ContextRunName = "__NAME__"
Variables ¶
var DefaultPubSubConfig = PubSubConfig{
CoProcess: 1,
ChannelBuffer: 128,
}
DefaultPubSubConfig 默认配置
Functions ¶
func GetContextRunId ¶ added in v1.0.5
GetContextRunId 获取run id
func GetContextRunName ¶ added in v1.0.5
GetContextRunName 获取run name
func SetContextRunId ¶ added in v1.0.5
SetContextRunId 设置run id
Types ¶
type ConProcessModelActuator ¶ added in v1.0.2
type ConProcessModelActuator struct {
// contains filtered or unexported fields
}
ConProcessModelActuator 并发处理模型 提供 Add 方法添加需要运行的函数,函数提供 context 变量,监听context的 Done() 确定是否结束处理 Start 启动方法执行,Stop停止方法执行 Wait 等到各个方法最终执行结束
func NewConProcessModelActuator ¶ added in v1.0.2
func NewConProcessModelActuator() *ConProcessModelActuator
func (*ConProcessModelActuator) Add ¶ added in v1.0.2
func (c *ConProcessModelActuator) Add(f CoFunc)
func (*ConProcessModelActuator) AddRepeat ¶ added in v1.0.2
func (c *ConProcessModelActuator) AddRepeat(f CoFunc, num int)
func (*ConProcessModelActuator) Start ¶ added in v1.0.2
func (c *ConProcessModelActuator) Start() error
func (*ConProcessModelActuator) Stop ¶ added in v1.0.2
func (c *ConProcessModelActuator) Stop() error
Stop Start运行后多次stop都可
func (*ConProcessModelActuator) Wait ¶ added in v1.0.2
func (c *ConProcessModelActuator) Wait() error
Wait 等待执行完成
type CtlRater ¶
type CtlRater struct {
StopNotifyWorker
Do func()
Rate int64
v1log.InvokeLog
// contains filtered or unexported fields
}
type EventBus ¶ added in v1.0.5
type EventBus = PubSubActuator[*extendparams.Message[interface{}]]
func NewEventBus ¶ added in v1.0.5
func NewEventBus(logger v1log.ILogger, opts ...PubSubOption) *EventBus
type FailAllFailCoMag ¶
type FailAllFailCoMag struct {
// contains filtered or unexported fields
}
FailAllFailCoMag 一次失败全部失败的CoMag 已废弃,请使用 ConProcessModelActuator 进行处理
func NewFailAllFailCoMag ¶
func NewFailAllFailCoMag() *FailAllFailCoMag
func (*FailAllFailCoMag) Add ¶
func (mg *FailAllFailCoMag) Add(f CoFunc)
func (*FailAllFailCoMag) Run ¶
func (mg *FailAllFailCoMag) Run()
type FragmentProcessor ¶ added in v1.0.5
type FragmentProcessor[T interface{}] struct {
// contains filtered or unexported fields
}
FragmentProcessor 分段自动处理,达到指定长度后才进行处理
func NewFragmentProcessor ¶ added in v1.0.5
func NewFragmentProcessor[T interface{}](fn func(...T) error, l int) *FragmentProcessor[T]
func (*FragmentProcessor[T]) AutoFinish ¶ added in v1.0.5
func (f *FragmentProcessor[T]) AutoFinish(t time.Duration)
AutoFinish 自动完成
func (*FragmentProcessor[T]) Finish ¶ added in v1.0.5
func (f *FragmentProcessor[T]) Finish() error
Finish 完成当前已有处理
func (*FragmentProcessor[T]) Push ¶ added in v1.0.5
func (f *FragmentProcessor[T]) Push(items ...T) error
Push 推入多个单位进行处理
func (*FragmentProcessor[T]) Start ¶ added in v1.0.5
func (f *FragmentProcessor[T]) Start() error
Start 开始异步处理
func (*FragmentProcessor[T]) Stop ¶ added in v1.0.5
func (f *FragmentProcessor[T]) Stop() error
Start 开始异步处理
type ProcessMultiSender ¶ added in v1.0.2
type ProcessMultiSender[T interface{}] struct {
// contains filtered or unexported fields
}
ProcessMultiSender 处理发送器(提供多个发送) 提供 Send 发送数据方法(与 ProcessSender 一致) 提供 Ch 获取收取数据的channel变量 与上面的 ProcessSender 的区别在于接收端可支持多个channel,可根据自身需要设置策略来动态分配对应的channel
func NewProcessMultiSender ¶ added in v1.0.2
func NewProcessMultiSender[T interface{}](num int) *ProcessMultiSender[T]
func (*ProcessMultiSender[T]) Ch ¶ added in v1.0.2
func (a *ProcessMultiSender[T]) Ch(index int) chan T
func (*ProcessMultiSender[T]) ChByData ¶ added in v1.0.2
func (a *ProcessMultiSender[T]) ChByData(d T) chan T
func (*ProcessMultiSender[T]) ChIndex ¶ added in v1.0.2
func (a *ProcessMultiSender[T]) ChIndex(d T) int
ChIndex 设置分发策略
func (*ProcessMultiSender[T]) Send ¶ added in v1.0.2
func (a *ProcessMultiSender[T]) Send(items ...T)
func (*ProcessMultiSender[T]) SetStrategy ¶ added in v1.0.2
func (a *ProcessMultiSender[T]) SetStrategy(f func(T) int)
SetStrategy 设置分发策略
type ProcessSender ¶ added in v1.0.2
type ProcessSender[T interface{}] struct {
// contains filtered or unexported fields
}
ProcessSender 处理发送器 提供 Send 发送数据方法 提供 Ch 获取收取数据的channel变量
func NewProcessSender ¶ added in v1.0.2
func NewProcessSender[T interface{}](l ...int) *ProcessSender[T]
func (*ProcessSender[T]) Ch ¶ added in v1.0.2
func (a *ProcessSender[T]) Ch() chan T
func (*ProcessSender[T]) Send ¶ added in v1.0.2
func (a *ProcessSender[T]) Send(items ...T)
type PubSubActuator ¶ added in v1.0.5
type PubSubActuator[T interface{}] struct {
appSupport.AppLifeSet
appSupport.NameSet
// contains filtered or unexported fields
}
PubSubActuator 事件总线(泛型)
func NewEventBusT ¶ added in v1.0.5
func NewEventBusT[T interface{}](logger v1log.ILogger, opts ...PubSubOption) *PubSubActuator[T]
func NewPubSubActuator ¶ added in v1.0.5
func NewPubSubActuator[T interface{}](logger v1log.ILogger, opts ...PubSubOption) *PubSubActuator[T]
func (*PubSubActuator[T]) Publish ¶ added in v1.0.5
func (bus *PubSubActuator[T]) Publish(topic string, msgs ...T) error
func (*PubSubActuator[T]) PublishByFn ¶ added in v1.0.5
func (bus *PubSubActuator[T]) PublishByFn(topic string, msgs ...T) error
PublishByFn 直接发布至fn执行(需要存在fn方法,同步执行)
func (*PubSubActuator[T]) Start ¶ added in v1.0.5
func (bus *PubSubActuator[T]) Start() error
func (*PubSubActuator[T]) Stop ¶ added in v1.0.5
func (bus *PubSubActuator[T]) Stop() error
func (*PubSubActuator[T]) StopAsync ¶ added in v1.0.5
func (bus *PubSubActuator[T]) StopAsync() error
func (*PubSubActuator[T]) Subscribe ¶ added in v1.0.5
func (bus *PubSubActuator[T]) Subscribe(topic string) (<-chan T, error)
func (*PubSubActuator[T]) SubscribeByFn ¶ added in v1.0.5
func (bus *PubSubActuator[T]) SubscribeByFn(topic string, fn func(int, T)) error
func (*PubSubActuator[T]) WaitStop ¶ added in v1.0.5
func (bus *PubSubActuator[T]) WaitStop()
type PubSubConfig ¶ added in v1.0.5
PubSubConfig 总线配置
type PubSubOption ¶ added in v1.0.5
type PubSubOption func(*PubSubConfig)
func WithPubSubChannelBuffer ¶ added in v1.0.5
func WithPubSubChannelBuffer(buffer int64) PubSubOption
func WithPubSubCoProcess ¶ added in v1.0.5
func WithPubSubCoProcess(num int) PubSubOption
type StopNotifyWorker ¶
type StopNotifyWorker struct {
StopWorker
// contains filtered or unexported fields
}
func (*StopNotifyWorker) Stop ¶
func (w *StopNotifyWorker) Stop()
func (*StopNotifyWorker) StopNotify ¶
func (w *StopNotifyWorker) StopNotify() chan bool
type StopWorker ¶
type StopWorker struct {
// contains filtered or unexported fields
}
func (*StopWorker) Start ¶
func (w *StopWorker) Start()
func (*StopWorker) Stop ¶
func (w *StopWorker) Stop()
func (*StopWorker) Stopped ¶
func (w *StopWorker) Stopped() bool