diff --git a/sundynix-shared/bus/bus.go b/sundynix-shared/bus/bus.go index 94cf78a..fc8fc4d 100644 --- a/sundynix-shared/bus/bus.go +++ b/sundynix-shared/bus/bus.go @@ -113,6 +113,7 @@ func (b *Bus) EnsureTaskStream(ctx context.Context) error { Name: contract.StreamTasks, Subjects: []string{contract.SubjectTasksAll}, Storage: jetstream.FileStorage, + Replicas: streamReplicas(), }) return err } @@ -341,6 +342,7 @@ func (b *Bus) EnsureStatusStream(ctx context.Context) error { Name: contract.StreamStatus, Subjects: []string{contract.SubjectTaskStatus}, Storage: jetstream.FileStorage, + Replicas: streamReplicas(), MaxAge: 24 * time.Hour, // 状态一天内必被消费,过期回收防无限增长 }) return err @@ -399,6 +401,7 @@ func (b *Bus) EnsureEvalStream(ctx context.Context) error { Name: contract.StreamEval, Subjects: []string{contract.SubjectEval}, Storage: jetstream.FileStorage, + Replicas: streamReplicas(), MaxAge: 24 * time.Hour, // 评测一天内必被消费,过期回收防无限增长 }) return err @@ -453,6 +456,7 @@ func (b *Bus) EnsureUsageStream(ctx context.Context) error { Name: contract.StreamUsage, Subjects: []string{contract.SubjectUsage}, Storage: jetstream.FileStorage, + Replicas: streamReplicas(), MaxAge: 72 * time.Hour, // 计费留足追赶窗口(网关长时间离线后仍能补齐) }) return err @@ -520,6 +524,7 @@ func (b *Bus) EnsureApprovalStream(ctx context.Context) error { Name: contract.StreamApprovals, Subjects: []string{contract.SubjectApprovalAll}, Storage: jetstream.FileStorage, + Replicas: streamReplicas(), MaxAge: 24 * time.Hour, }) return err @@ -774,6 +779,18 @@ func DrainTimeout() time.Duration { return 30 * time.Second } +// streamReplicas 是 JetStream 流/KV 桶的副本数,经 NATS_STREAM_REPLICAS 配置,缺省 1。 +// 单机(单节点 NATS)保持 1;三机/单机多节点集群部署设 3,durable 流才有 Raft quorum 容错 +// (计费/状态/用量/评测/审批/入库回写不会因单节点挂掉而丢)。合法值 1/3/5(Raft 需奇数)。 +func streamReplicas() int { + if v := os.Getenv("NATS_STREAM_REPLICAS"); v != "" { + if n, err := strconv.Atoi(v); err == nil && (n == 1 || n == 3 || n == 5) { + return n + } + } + return 1 +} + // drainWait 停止接新活后,等待在途 WaitGroup 跑完,或 dctx 到期放弃。 func drainWait(wg *sync.WaitGroup, dctx context.Context) { done := make(chan struct{}) @@ -890,6 +907,7 @@ func (b *Bus) EnsureIngestStream(ctx context.Context) error { Name: contract.StreamIngest, Subjects: []string{contract.SubjectIngestAll}, Storage: jetstream.FileStorage, + Replicas: streamReplicas(), MaxAge: ingestStreamMaxAge, }) return err @@ -934,7 +952,7 @@ func (b *Bus) ConsumeIngestJobs(ctx context.Context, h IngestHandler) (drain fun // 一篇几十万字含切块+上千次 embedding+至多 60 次图谱抽取,可达数分钟; // AckWait 须覆盖单篇最长入库时长,否则在途未 ack 会被重投成重复入库。 AckWait: ingestAckWait(), - MaxAckPending: concurrency, // 背压:服务端不下发超过本节点同时能处理的量 + MaxAckPending: concurrency, // 背压:服务端不下发超过本节点同时能处理的量 MaxDeliver: ingestMaxDeliver, // 毒消息兜底:重投上限,避免坏作业无限循环 }) if err != nil { @@ -1002,9 +1020,10 @@ type KVHandle struct{ kv jetstream.KeyValue } // 决定到达即恢复;TTL 兜底清理被遗弃的中断(如审批人始终不处理)。 func (b *Bus) Checkpoints(ctx context.Context, bucket string, ttl time.Duration) (*KVHandle, error) { kv, err := b.js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{ - Bucket: bucket, - Storage: jetstream.FileStorage, - TTL: ttl, + Bucket: bucket, + Storage: jetstream.FileStorage, + Replicas: streamReplicas(), + TTL: ttl, }) if err != nil { return nil, fmt.Errorf("kv bucket %q: %w", bucket, err)