comag

package
v1.0.5 Latest Latest
Warning

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

Go to latest
Published: Aug 27, 2026 License: Apache-2.0 Imports: 11 Imported by: 0

README

并发处理功能

基本使用

消息发送器 - ProcessSender

对应模型类: ProcessSender

基本使用
sender := NewProcessSender[*tests.Person]()

go func(ctx context.Context) {
    for {
        select {
        case d := <-sender.Ch():
            t.Logf("[%d-%s]%v", ctx.Value(ContextRunId), ctx.Value(ContextRunName), d)
        case <-ctx.Done():
            return
        }
    }
}

for i := 1; i <= 100; i++ {
    sender.Send(tests.NewPerson())
    time.Sleep(time.Millisecond * 200)
}
并发处理器 - ConProcessModelActuator

对应模型类: ConProcessModelActuator

基本使用
actuator := NewConProcessModelActuator()
f := func(id int) CoFunc {
    return func(ctx context.Context) {
        defer fmt.Printf("xxx-%d DONE\n", id)
        ti := time.NewTicker(time.Second * time.Duration(id%4))
        for {
            select {
            case <-ctx.Done():
                time.Sleep(time.Second * time.Duration(id%4))
                return
            case <-ti.C:
                fmt.Printf("DOING-%d\n", id)
                if id == 3 {
                    a.Stop()
                    return
                }
            }
        }
    }
}
actuator.Add(f(1))
actuator.Add(f(2))
actuator.Add(f(3))
actuator.Start()

停止执行(阻塞等待结束)
time.Sleep(time.Second * 3)
actuator.Stop()
fmt.Printf("TO END\n")

actuator.Wait()

a.Stop()
fmt.Printf("TO END\n")

配合ProcessSender(并发消息处理)

可配合 ProcessSender 进行并发消息的分配处理

sender := NewProcessSender[int]()
actuator := NewConProcessModelActuator()

f := func(ctx context.Context) {
		for {
			select {
			case d := <-sender.Ch():
				t.Logf("[%d-%s]%d", ctx.Value(ContextRunId), ctx.Value(ContextRunName), d)
			case <-ctx.Done():
				return
			}
		}
}
actuator.AddRepeat(f, 10)
actuator.Start()
	
for i := 1; i <= 100; i++ {
		sender.Send(i)
		time.Sleep(time.Millisecond * 200)
}
发布订阅处理器 - PubSubActuator

对应模型类: PubSubActuator

初始化订阅
bus := NewPubSubActuator[int](v1log.Out)

countA := 0
common.ErrPanic(bus.SubscribeByFn("order.created", func(index int, msg int) {
    countA++
}))

common.ErrPanic(bus.Start())

手动启动
common.ErrPanic(bus.Start())
发送消息
common.ErrPanic(bus.Publish("order.created", 1))
common.ErrPanic(bus.Publish("order.created", 2))

手动停止(阻塞等待结束)
common.ErrPanic(bus.Stop())
支持泛型消息

指定int类型作为消息数据

bus := NewPubSubActuator[int](v1log.Out)



使用自定义结构体作为消息数据


bus := NewPubSubActuator[*tests.Person](v1log.Out)

EventBus(扩展功能)

支持事件总线,固定消息数据为 Message 类型

bus := NewEventBus(v1log.Out)
err = bus.SubscribeByFn("order.created", func(index int, msg *extendparams.Message[interface{}]) {
	countA++
	t.Logf("[GET.created-%d]%s %d %+v", index, msg.ExtendParams.Get("topic"), countA, msg.PayloadData)
})
common.ErrPanic(err)

common.ErrPanic(bus.Start())

msg := extendparams.NewMessage[interface{}]()
msg.PayloadData = tests.NewPerson()
msg.ExtendParams.Set("topic", "order.payed")
msg.ExtendParams.Set("index", fmt.Sprint(i))
common.ErrPanic(bus.Publish("order.payed", msg))

Documentation

Index

Constants

View Source
const ContextRunId = "__ID__"
View Source
const ContextRunName = "__NAME__"

Variables

View Source
var DefaultPubSubConfig = PubSubConfig{
	CoProcess:     1,
	ChannelBuffer: 128,
}

DefaultPubSubConfig 默认配置

Functions

func DemoCoA

func DemoCoA(ctx context.Context)

func DemoCoB

func DemoCoB(ctx context.Context)

func GetContextRunId added in v1.0.5

func GetContextRunId(ctx context.Context) int

GetContextRunId 获取run id

func GetContextRunName added in v1.0.5

func GetContextRunName(ctx context.Context) string

GetContextRunName 获取run name

func SetContextRunId added in v1.0.5

func SetContextRunId(ctx context.Context, index int) context.Context

SetContextRunId 设置run id

func SetContextRunName added in v1.0.5

func SetContextRunName(ctx context.Context, index string) context.Context

SetContextRunName 设置run name

Types

type CoFunc

type CoFunc func(context.Context)

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
}

func (*CtlRater) Start

func (ctl *CtlRater) Start()

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

type PubSubConfig struct {
	CoProcess     int   //并发处理数
	ChannelBuffer int64 //发送通道缓存大小(对应channel size)
}

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

Jump to

Keyboard shortcuts

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