diff --git a/docs/architecture/iot-data-platform-principles.md b/docs/architecture/iot-data-platform-principles.md index e6d11e01..ea90d4c2 100644 --- a/docs/architecture/iot-data-platform-principles.md +++ b/docs/architecture/iot-data-platform-principles.md @@ -71,7 +71,7 @@ Every received frame should become one `FrameEnvelope`. 1. Stabilize ingress: connection lifecycle, protocol parser correctness, bounded backpressure, durable publish. 2. Stabilize event log: NATS to Kafka bridge, topic names, partition keys, retry and replay. 3. Simplify storage: keep only the minimal tables above, remove duplicated payloads from business tables. -4. Stabilize realtime: Redis full latest raw state, MySQL lightweight snapshot and location. +4. Stabilize realtime: Redis full latest raw state, MySQL lightweight snapshot and location. Realtime consumers start from latest when no committed offset exists; historical rebuilds must be explicit replay jobs, not accidental backlog scans. 5. Stabilize history: raw and location query pagination, clear time zone behavior. 6. Stabilize metrics: idempotent daily metrics from mileage differences and raw replay. 7. Add operations: health, readiness, metrics, lag, connection counts, deploy notes. diff --git a/docs/ops/vehicle-ingest-runbook.md b/docs/ops/vehicle-ingest-runbook.md index bb1550a9..b37204ea 100644 --- a/docs/ops/vehicle-ingest-runbook.md +++ b/docs/ops/vehicle-ingest-runbook.md @@ -38,6 +38,8 @@ 如果 stat-writer 少消费某个 raw topic,对应协议的每日里程不会进入 `vehicle_daily_mileage`。修正后可能出现短时间 Kafka lag,这是在追补历史 backlog;只要 `vehicle_stat_writes_total` 持续增长且 lag 下降,就是健康状态。 +Realtime API 是当前态投影,默认在没有已提交 offset 时从 latest 开始消费。切换 topic 或新建 consumer group 后,不应让 realtime 追扫历史 raw backlog;需要重建当前态时,应使用明确的回放任务或手动 reset offset。 + 业务端口: - GB32960 TCP:`0.0.0.0:32960` diff --git a/go/vehicle-gateway/cmd/realtime-api/main.go b/go/vehicle-gateway/cmd/realtime-api/main.go index 7701207a..228a6132 100644 --- a/go/vehicle-gateway/cmd/realtime-api/main.go +++ b/go/vehicle-gateway/cmd/realtime-api/main.go @@ -167,13 +167,7 @@ func consumeKafka(ctx context.Context, logger interface { Warn(string, ...any) }, registry *metrics.Registry, updater realtimeUpdater, brokers []string) { kafkaTopics := kafkaTopicsFromEnv() - reader := kafka.NewReader(kafka.ReaderConfig{ - Brokers: brokers, - GroupID: env("KAFKA_GROUP", "go-realtime-api"), - GroupTopics: kafkaTopics, - MinBytes: 1, - MaxBytes: 10e6, - }) + reader := kafka.NewReader(realtimeReaderConfig(brokers, kafkaTopics)) defer reader.Close() logger.Info("realtime kafka consumer started", "topics", strings.Join(kafkaTopics, ",")) for { @@ -189,6 +183,17 @@ func consumeKafka(ctx context.Context, logger interface { } } +func realtimeReaderConfig(brokers []string, kafkaTopics []string) kafka.ReaderConfig { + return kafka.ReaderConfig{ + Brokers: brokers, + GroupID: env("KAFKA_GROUP", "go-realtime-api"), + GroupTopics: kafkaTopics, + MinBytes: 1, + MaxBytes: 10e6, + StartOffset: kafka.LastOffset, + } +} + func kafkaTopicsFromEnv() []string { defaultTopics := strings.Join([]string{topics.RawGB32960, topics.RawJT808, topics.RawYutongMQTT}, ",") return splitCSV(env("KAFKA_TOPICS", defaultTopics)) diff --git a/go/vehicle-gateway/cmd/realtime-api/main_test.go b/go/vehicle-gateway/cmd/realtime-api/main_test.go index 0c242924..7fa4c491 100644 --- a/go/vehicle-gateway/cmd/realtime-api/main_test.go +++ b/go/vehicle-gateway/cmd/realtime-api/main_test.go @@ -104,6 +104,17 @@ func TestKafkaTopicsFromEnvDefaultsToGoRawTopics(t *testing.T) { } } +func TestRealtimeReaderConfigStartsAtLatestOffset(t *testing.T) { + cfg := realtimeReaderConfig([]string{"127.0.0.1:9092"}, []string{"vehicle.raw.go.jt808.v1"}) + + if cfg.StartOffset != kafka.LastOffset { + t.Fatalf("StartOffset = %d, want kafka.LastOffset %d", cfg.StartOffset, kafka.LastOffset) + } + if got := strings.Join(cfg.GroupTopics, ","); got != "vehicle.raw.go.jt808.v1" { + t.Fatalf("GroupTopics = %q", got) + } +} + func TestRealtimeAPIDoesNotExposeDuplicateMileagePointRoute(t *testing.T) { source, err := os.ReadFile("main.go") if err != nil {