1
This commit is contained in:
@@ -91,3 +91,44 @@ func (p *LabelTaskPool) Submit(ctx context.Context, fn func(ctx context.Context)
|
||||
return ctx.Err()
|
||||
}
|
||||
}
|
||||
|
||||
// GenTaskPool 文生图任务并发池:逐张调 local-ai 生成图片(IO 等待为主),
|
||||
// 并发度来自 config.yml imageGen.poolSize(z-image 显存独占,默认 1),
|
||||
// 缺失或非法时回退 consts 默认值。池内任务禁止提交本池(防 worker 饿死死锁)。
|
||||
type GenTaskPool struct {
|
||||
pool *grpool.Pool
|
||||
}
|
||||
|
||||
var (
|
||||
genTaskPoolOnce sync.Once
|
||||
genTaskPool *GenTaskPool
|
||||
)
|
||||
|
||||
// GenTaskPoolInstance 进程级文生图池单例(懒初始化,读取配置)。
|
||||
func GenTaskPoolInstance() *GenTaskPool {
|
||||
genTaskPoolOnce.Do(func() {
|
||||
ctx := context.Background()
|
||||
size := g.Cfg().MustGet(ctx, "imageGen.poolSize", consts.GenPoolDefaultSize).Int()
|
||||
if size <= 0 {
|
||||
size = consts.GenPoolDefaultSize
|
||||
}
|
||||
genTaskPool = &GenTaskPool{pool: grpool.New(size, size)}
|
||||
})
|
||||
return genTaskPool
|
||||
}
|
||||
|
||||
// Submit 提交单张图片的生成任务并等待完成,返回任务的 error。
|
||||
func (p *GenTaskPool) Submit(ctx context.Context, fn func(ctx context.Context) error) error {
|
||||
res := make(chan error, 1)
|
||||
if err := p.pool.Add(ctx, func(ctx context.Context) {
|
||||
res <- fn(ctx)
|
||||
}); err != nil {
|
||||
return err
|
||||
}
|
||||
select {
|
||||
case err := <-res:
|
||||
return err
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user