feat(go): batch nats fast writer tdengine writes
This commit is contained in:
@@ -39,7 +39,7 @@
|
|||||||
- Gateway 默认 `TCP_MAX_CONNECTIONS` 已调整为 `120000`,生产仍可通过环境变量覆盖。
|
- Gateway 默认 `TCP_MAX_CONNECTIONS` 已调整为 `120000`,生产仍可通过环境变量覆盖。
|
||||||
- 已新增可重复的 TCP 连接压测工具,后续需要用它跑 10K/50K/100K 阶段测试并记录结果。
|
- 已新增可重复的 TCP 连接压测工具,后续需要用它跑 10K/50K/100K 阶段测试并记录结果。
|
||||||
- 仍缺读超时按协议维度计数;Gateway frame duration 已提供 histogram,可用于 parse + enqueue + response 的 p95/p99 估算。
|
- 仍缺读超时按协议维度计数;Gateway frame duration 已提供 histogram,可用于 parse + enqueue + response 的 p95/p99 估算。
|
||||||
- TDengine history writer 当前逐消息插入,后续需要批量写入。
|
- TDengine history writer 和 NATS fast writer 已使用 batch 写入;后续需要用带帧率压测验证 batch size、flush latency 和存储端承载能力。
|
||||||
- Kafka topic 当前 12 分区,100K 目标下需要结合实际 FPS 再评估分区数。
|
- Kafka topic 当前 12 分区,100K 目标下需要结合实际 FPS 再评估分区数。
|
||||||
|
|
||||||
## ECS OS 参数建议
|
## ECS OS 参数建议
|
||||||
|
|||||||
@@ -421,6 +421,8 @@ vehicle_fast_writer_stage_duration_ms_histogram_sum{subject,stage,status}
|
|||||||
|
|
||||||
`stage` 取值包括 `tdengine`、`redis`、`ack`,分别对应快速写 TDengine raw/location、Redis realtime 投影、NATS ack。带帧率压测时如果某个 stage 的 p99 持续上升,可以直接定位 fast path 是存储瓶颈还是 ack/网络瓶颈。
|
`stage` 取值包括 `tdengine`、`redis`、`ack`,分别对应快速写 TDengine raw/location、Redis realtime 投影、NATS ack。带帧率压测时如果某个 stage 的 p99 持续上升,可以直接定位 fast path 是存储瓶颈还是 ack/网络瓶颈。
|
||||||
|
|
||||||
|
2026-07-03 后续优化:`nats-fast-writer` 已从“Fetch batch 后逐条写 TDengine”改为“Fetch batch 后先调用 `history.Writer.AppendAllBatch` 批量写 TDengine,再逐条更新 Redis 和 ack”。这样可以减少 TDengine round trip;Redis/ack 仍按单条执行,避免某条实时投影失败时误 ack。
|
||||||
|
|
||||||
### NATS Kafka bridge batch 指标
|
### NATS Kafka bridge batch 指标
|
||||||
|
|
||||||
2026-07-03 已为 NATS -> Kafka bridge 增加 batch pending 和 duration histogram:
|
2026-07-03 已为 NATS -> Kafka bridge 增加 batch pending 和 duration histogram:
|
||||||
|
|||||||
@@ -184,6 +184,7 @@ func loadConfig() config {
|
|||||||
|
|
||||||
type fastAppender interface {
|
type fastAppender interface {
|
||||||
AppendAll(context.Context, envelope.FrameEnvelope) error
|
AppendAll(context.Context, envelope.FrameEnvelope) error
|
||||||
|
AppendAllBatch(context.Context, []envelope.FrameEnvelope) error
|
||||||
}
|
}
|
||||||
|
|
||||||
type fastUpdater interface {
|
type fastUpdater interface {
|
||||||
@@ -210,23 +211,27 @@ func runFastWorker(ctx context.Context, logger *slog.Logger, registry *metrics.R
|
|||||||
time.Sleep(time.Second)
|
time.Sleep(time.Second)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
fastMessages := make([]*fastMessage, 0, len(msgs))
|
||||||
for _, msg := range msgs {
|
for _, msg := range msgs {
|
||||||
msg := msg
|
natsMsg := msg
|
||||||
operationCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), cfg.OperationWait)
|
fastMessages = append(fastMessages, &fastMessage{
|
||||||
err := processFastMessage(operationCtx, registry, appender, updater, &fastMessage{
|
subject: natsMsg.Subject,
|
||||||
subject: msg.Subject,
|
data: natsMsg.Data,
|
||||||
data: msg.Data,
|
|
||||||
ack: func() error {
|
ack: func() error {
|
||||||
return msg.Ack()
|
return natsMsg.Ack()
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
cancel()
|
}
|
||||||
if err != nil {
|
operationCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), cfg.OperationWait)
|
||||||
addFastMetric(registry, msg.Subject, "error")
|
err = processFastBatch(operationCtx, registry, appender, updater, fastMessages)
|
||||||
logger.Error("fast write failed", "subject", msg.Subject, "error", err)
|
cancel()
|
||||||
continue
|
if err != nil {
|
||||||
}
|
addFastMetric(registry, fastBatchSubject(fastMessages), "error")
|
||||||
addFastMetric(registry, msg.Subject, "ok")
|
logger.Error("fast write batch failed", "messages", len(fastMessages), "error", err)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
for _, msg := range fastMessages {
|
||||||
|
addFastMetric(registry, msg.subject, "ok")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -237,6 +242,53 @@ type natsPullSubscription interface {
|
|||||||
|
|
||||||
var fastWriterStageDurationBucketsMS = []float64{1, 5, 10, 25, 50, 100, 250, 500, 1000, 5000}
|
var fastWriterStageDurationBucketsMS = []float64{1, 5, 10, 25, 50, 100, 250, 500, 1000, 5000}
|
||||||
|
|
||||||
|
func processFastBatch(ctx context.Context, registry *metrics.Registry, appender fastAppender, updater fastUpdater, messages []*fastMessage) error {
|
||||||
|
if len(messages) == 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
envelopes := make([]envelope.FrameEnvelope, 0, len(messages))
|
||||||
|
validMessages := make([]*fastMessage, 0, len(messages))
|
||||||
|
for _, msg := range messages {
|
||||||
|
var env envelope.FrameEnvelope
|
||||||
|
if err := json.Unmarshal(msg.data, &env); err != nil {
|
||||||
|
if msg.ack != nil {
|
||||||
|
_ = msg.ack()
|
||||||
|
}
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
envelopes = append(envelopes, env)
|
||||||
|
validMessages = append(validMessages, msg)
|
||||||
|
}
|
||||||
|
if len(envelopes) == 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
subject := fastBatchSubject(validMessages)
|
||||||
|
started := time.Now()
|
||||||
|
err := appender.AppendAllBatch(ctx, envelopes)
|
||||||
|
recordFastWriterStageDuration(registry, subject, "tdengine", statusFromError(err), time.Since(started))
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("tdengine batch append: %w", err)
|
||||||
|
}
|
||||||
|
for i, env := range envelopes {
|
||||||
|
msg := validMessages[i]
|
||||||
|
started = time.Now()
|
||||||
|
err = updater.FastUpdate(ctx, env)
|
||||||
|
recordFastWriterStageDuration(registry, msg.subject, "redis", statusFromError(err), time.Since(started))
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("redis fast update: %w", err)
|
||||||
|
}
|
||||||
|
if msg.ack != nil {
|
||||||
|
started = time.Now()
|
||||||
|
err = msg.ack()
|
||||||
|
recordFastWriterStageDuration(registry, msg.subject, "ack", statusFromError(err), time.Since(started))
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("nats ack: %w", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
func processFastMessage(ctx context.Context, registry *metrics.Registry, appender fastAppender, updater fastUpdater, msg *fastMessage) error {
|
func processFastMessage(ctx context.Context, registry *metrics.Registry, appender fastAppender, updater fastUpdater, msg *fastMessage) error {
|
||||||
var env envelope.FrameEnvelope
|
var env envelope.FrameEnvelope
|
||||||
if err := json.Unmarshal(msg.data, &env); err != nil {
|
if err := json.Unmarshal(msg.data, &env); err != nil {
|
||||||
@@ -268,6 +320,22 @@ func processFastMessage(ctx context.Context, registry *metrics.Registry, appende
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func fastBatchSubject(messages []*fastMessage) string {
|
||||||
|
if len(messages) == 0 {
|
||||||
|
return "unknown"
|
||||||
|
}
|
||||||
|
subject := messages[0].subject
|
||||||
|
for _, msg := range messages[1:] {
|
||||||
|
if msg.subject != subject {
|
||||||
|
return "mixed"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if strings.TrimSpace(subject) == "" {
|
||||||
|
return "unknown"
|
||||||
|
}
|
||||||
|
return subject
|
||||||
|
}
|
||||||
|
|
||||||
func ensureStream(js nats.JetStreamContext, cfg config) error {
|
func ensureStream(js nats.JetStreamContext, cfg config) error {
|
||||||
stream := &nats.StreamConfig{
|
stream := &nats.StreamConfig{
|
||||||
Name: cfg.NATSStream,
|
Name: cfg.NATSStream,
|
||||||
|
|||||||
@@ -80,9 +80,44 @@ func TestProcessFastMessageRecordsStageDurationMetrics(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestProcessFastBatchAppendsTDengineBatchBeforeRedisAndAck(t *testing.T) {
|
||||||
|
first := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, VIN: "VIN001", EventID: "evt-4"}
|
||||||
|
second := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, VIN: "VIN002", EventID: "evt-5"}
|
||||||
|
firstPayload, err := first.MarshalJSONBytes()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
secondPayload, err := second.MarshalJSONBytes()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
appender := &recordingFastAppender{}
|
||||||
|
updater := &recordingFastUpdater{}
|
||||||
|
ackCount := 0
|
||||||
|
msgs := []*fastMessage{
|
||||||
|
{subject: "vehicle.raw.go.jt808.v1", data: firstPayload, ack: func() error { ackCount++; return nil }},
|
||||||
|
{subject: "vehicle.raw.go.jt808.v1", data: secondPayload, ack: func() error { ackCount++; return nil }},
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := processFastBatch(context.Background(), nil, appender, updater, msgs); err != nil {
|
||||||
|
t.Fatalf("processFastBatch() error = %v", err)
|
||||||
|
}
|
||||||
|
if appender.count != 0 {
|
||||||
|
t.Fatalf("AppendAll count = %d, want 0", appender.count)
|
||||||
|
}
|
||||||
|
if appender.batchCount != 1 || appender.batchRows != 2 {
|
||||||
|
t.Fatalf("AppendAllBatch count=%d rows=%d, want count=1 rows=2", appender.batchCount, appender.batchRows)
|
||||||
|
}
|
||||||
|
if updater.count != 2 || ackCount != 2 {
|
||||||
|
t.Fatalf("updates=%d acks=%d, want 2/2", updater.count, ackCount)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
type recordingFastAppender struct {
|
type recordingFastAppender struct {
|
||||||
count int
|
count int
|
||||||
err error
|
batchCount int
|
||||||
|
batchRows int
|
||||||
|
err error
|
||||||
}
|
}
|
||||||
|
|
||||||
func (a *recordingFastAppender) AppendAll(context.Context, envelope.FrameEnvelope) error {
|
func (a *recordingFastAppender) AppendAll(context.Context, envelope.FrameEnvelope) error {
|
||||||
@@ -90,6 +125,12 @@ func (a *recordingFastAppender) AppendAll(context.Context, envelope.FrameEnvelop
|
|||||||
return a.err
|
return a.err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (a *recordingFastAppender) AppendAllBatch(_ context.Context, envs []envelope.FrameEnvelope) error {
|
||||||
|
a.batchCount++
|
||||||
|
a.batchRows += len(envs)
|
||||||
|
return a.err
|
||||||
|
}
|
||||||
|
|
||||||
type recordingFastUpdater struct {
|
type recordingFastUpdater struct {
|
||||||
count int
|
count int
|
||||||
err error
|
err error
|
||||||
|
|||||||
Reference in New Issue
Block a user