perf(dispatcher): 任务并发消费 —— 待审/慢任务不再阻塞其他任务
bus.ConsumeTasks 由单条串行改为限并发分发:每个任务进独立 worker goroutine,
并发上限 = 信号量 + 消费者 MaxAckPending(DISPATCHER_CONCURRENCY,默认 8)。
- 背压:并发满则在 select{sem, ctx.Done} 处等空位;关停时留消息不 ack 待重投。
- 健壮:worker 内 recover panic → Term(避免崩溃循环);ack/nak/span 收口在 worker。
- 并发安全已核:CircuitBreaker 有锁、Orchestrator.turns 有 turnMu、pool RWMutex、
evaluate 本就异步。
效果:一个 HITL 待审任务(Handle 阻塞至多 5min)或长 LLM 生成不再冻结后续任务。
验证:单测 TestConcurrentConsume(A 阻塞时 B 完成);live 实测 HITL 停 waiting 期间
普通任务 3s 跑完且 HITL 不受影响。全模块 build+vet+test 全绿。
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
+45
-14
@@ -6,6 +6,9 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"log"
|
||||||
|
"os"
|
||||||
|
"strconv"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/nats-io/nats.go"
|
"github.com/nats-io/nats.go"
|
||||||
@@ -441,9 +444,21 @@ func (b *Bus) SubscribeConfigUpdated(kind string, onUpdate func(*contract.ModelC
|
|||||||
// TaskHandler 处理一个消费到的任务。
|
// TaskHandler 处理一个消费到的任务。
|
||||||
type TaskHandler func(ctx context.Context, t *contract.Task) error
|
type TaskHandler func(ctx context.Context, t *contract.Task) error
|
||||||
|
|
||||||
|
// taskConcurrency 返回任务并发处理上限(env DISPATCHER_CONCURRENCY,默认 8)。
|
||||||
|
func taskConcurrency() int {
|
||||||
|
if v := os.Getenv("DISPATCHER_CONCURRENCY"); v != "" {
|
||||||
|
if n, err := strconv.Atoi(v); err == nil && n > 0 {
|
||||||
|
return n
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return 8
|
||||||
|
}
|
||||||
|
|
||||||
// ConsumeTasks 在持久消费者上消费任务,队列组内负载均衡。
|
// ConsumeTasks 在持久消费者上消费任务,队列组内负载均衡。
|
||||||
// 返回的 stop 函数用于优雅停止消费。
|
// 每个任务分发到独立 worker goroutine 并发执行——一个慢任务/HITL 待审不再阻塞后续任务。
|
||||||
|
// 并发上限由信号量 + 消费者 MaxAckPending 双重约束(背压)。返回的 stop 用于优雅停止消费。
|
||||||
func (b *Bus) ConsumeTasks(ctx context.Context, h TaskHandler) (stop func(), err error) {
|
func (b *Bus) ConsumeTasks(ctx context.Context, h TaskHandler) (stop func(), err error) {
|
||||||
|
concurrency := taskConcurrency()
|
||||||
cons, err := b.js.CreateOrUpdateConsumer(ctx, contract.StreamTasks, jetstream.ConsumerConfig{
|
cons, err := b.js.CreateOrUpdateConsumer(ctx, contract.StreamTasks, jetstream.ConsumerConfig{
|
||||||
Durable: contract.ConsumerDurable,
|
Durable: contract.ConsumerDurable,
|
||||||
AckPolicy: jetstream.AckExplicitPolicy,
|
AckPolicy: jetstream.AckExplicitPolicy,
|
||||||
@@ -451,31 +466,47 @@ func (b *Bus) ConsumeTasks(ctx context.Context, h TaskHandler) (stop func(), err
|
|||||||
// HITL:审批节点会让 Handle 阻塞等人工决定(最长约 5 分钟),
|
// HITL:审批节点会让 Handle 阻塞等人工决定(最长约 5 分钟),
|
||||||
// AckWait 必须覆盖「审批等待 + 图执行」总时长,否则消息在途未 ack 会被重投成重复任务。
|
// AckWait 必须覆盖「审批等待 + 图执行」总时长,否则消息在途未 ack 会被重投成重复任务。
|
||||||
AckWait: 15 * time.Minute,
|
AckWait: 15 * time.Minute,
|
||||||
|
// 在途未 ack 上限 = 并发上限:服务端不会下发超过本节点同时能处理的量(背压)。
|
||||||
|
MaxAckPending: concurrency,
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("create consumer: %w", err)
|
return nil, fmt.Errorf("create consumer: %w", err)
|
||||||
}
|
}
|
||||||
|
sem := make(chan struct{}, concurrency) // 限并发:最多 N 个任务同时执行
|
||||||
cc, err := cons.Consume(func(msg jetstream.Msg) {
|
cc, err := cons.Consume(func(msg jetstream.Msg) {
|
||||||
t, err := contract.Unmarshal(msg.Data())
|
t, err := contract.Unmarshal(msg.Data())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
_ = msg.Term() // 脏数据,丢弃不重投
|
_ = msg.Term() // 脏数据,丢弃不重投
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
// 从消息头还原上游链路,开消费 span(成为 gateway 发布 span 的子节点)。
|
// 背压:并发已满则在此等空位;正在关停则留消息不 ack(稍后重投)。
|
||||||
mctx := extractTrace(ctx, nats.Header(msg.Headers()))
|
select {
|
||||||
mctx, span := tracer().Start(mctx, "nats.consume task",
|
case sem <- struct{}{}:
|
||||||
trace.WithSpanKind(trace.SpanKindConsumer),
|
case <-ctx.Done():
|
||||||
trace.WithAttributes(attribute.String("sundynix.task_id", t.ID)))
|
|
||||||
herr := h(mctx, t)
|
|
||||||
if herr != nil {
|
|
||||||
span.RecordError(herr)
|
|
||||||
}
|
|
||||||
span.End()
|
|
||||||
if herr != nil {
|
|
||||||
_ = msg.NakWithDelay(time.Second) // 处理失败,延迟重投
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
_ = msg.Ack()
|
go func() {
|
||||||
|
defer func() {
|
||||||
|
<-sem // 释放并发额度
|
||||||
|
if r := recover(); r != nil {
|
||||||
|
// 任务处理 panic:丢弃不重投(避免崩溃循环),记录后继续。
|
||||||
|
log.Printf("[bus] task %s handler panic: %v", t.ID, r)
|
||||||
|
_ = msg.Term()
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
// 从消息头还原上游链路,开消费 span(成为 gateway 发布 span 的子节点)。
|
||||||
|
mctx := extractTrace(ctx, nats.Header(msg.Headers()))
|
||||||
|
mctx, span := tracer().Start(mctx, "nats.consume task",
|
||||||
|
trace.WithSpanKind(trace.SpanKindConsumer),
|
||||||
|
trace.WithAttributes(attribute.String("sundynix.task_id", t.ID)))
|
||||||
|
defer span.End()
|
||||||
|
if herr := h(mctx, t); herr != nil {
|
||||||
|
span.RecordError(herr)
|
||||||
|
_ = msg.NakWithDelay(time.Second) // 处理失败,延迟重投
|
||||||
|
return
|
||||||
|
}
|
||||||
|
_ = msg.Ack()
|
||||||
|
}()
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("consume: %w", err)
|
return nil, fmt.Errorf("consume: %w", err)
|
||||||
|
|||||||
@@ -194,3 +194,63 @@ func TestTokenStreamRoundTrip(t *testing.T) {
|
|||||||
t.Fatal("timeout: 未收到流结束信号")
|
t.Fatal("timeout: 未收到流结束信号")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestConcurrentConsume 验证并发消费:一个长时间阻塞的任务(模拟 HITL 待审)
|
||||||
|
// 不应阻塞后续任务——旧的串行消费下本测试会超时失败。
|
||||||
|
func TestConcurrentConsume(t *testing.T) {
|
||||||
|
url := startEmbeddedNATS(t)
|
||||||
|
|
||||||
|
gw, err := bus.Connect(url)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("gateway connect: %v", err)
|
||||||
|
}
|
||||||
|
defer gw.Close()
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
if err := gw.EnsureTaskStream(ctx); err != nil {
|
||||||
|
t.Fatalf("ensure stream: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
dp, err := bus.Connect(url)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("dispatcher connect: %v", err)
|
||||||
|
}
|
||||||
|
defer dp.Close()
|
||||||
|
|
||||||
|
blockA := make(chan struct{})
|
||||||
|
doneB := make(chan string, 1)
|
||||||
|
stop, err := dp.ConsumeTasks(ctx, func(_ context.Context, task *contract.Task) error {
|
||||||
|
switch task.ID {
|
||||||
|
case "A":
|
||||||
|
<-blockA // 模拟审批长阻塞
|
||||||
|
case "B":
|
||||||
|
doneB <- task.ID
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("consume: %v", err)
|
||||||
|
}
|
||||||
|
defer stop()
|
||||||
|
defer close(blockA)
|
||||||
|
|
||||||
|
// 先发 A 并等它被消费、卡在 handler 里;再发 B。
|
||||||
|
if _, err := gw.PublishTask(ctx, &contract.Task{ID: "A", Graph: json.RawMessage(`{}`)}); err != nil {
|
||||||
|
t.Fatalf("publish A: %v", err)
|
||||||
|
}
|
||||||
|
time.Sleep(400 * time.Millisecond)
|
||||||
|
if _, err := gw.PublishTask(ctx, &contract.Task{ID: "B", Graph: json.RawMessage(`{}`)}); err != nil {
|
||||||
|
t.Fatalf("publish B: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
select {
|
||||||
|
case id := <-doneB:
|
||||||
|
if id != "B" {
|
||||||
|
t.Fatalf("unexpected task done: %q", id)
|
||||||
|
}
|
||||||
|
t.Log("✓ 并发消费生效:A 仍阻塞时 B 已完成")
|
||||||
|
case <-time.After(4 * time.Second):
|
||||||
|
t.Fatal("B 被阻塞的 A 卡住——并发消费未生效(仍是串行)")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user