From e80fa6f38cfc8df0f220344c02f0ebdc71884d88 Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 10:56:33 +0800 Subject: [PATCH] refactor: rename kafka sink module --- CHANGELOG.md | 2 +- README.md | 2 +- docs/module-data-flow.html | 6 +-- .../gb32960-service-split-runbook.md | 2 +- ...-20-gb32960-body-parser-block-isolation.md | 2 +- .../plans/2026-06-23-gb32960-service-split.md | 46 +++++++++---------- ...26-06-29-vehicle-ingest-redesign-phase1.md | 38 +++++++-------- ...2026-06-23-gb32960-service-split-design.md | 6 +-- modules/apps/gb32960-ingest-app/pom.xml | 2 +- modules/apps/jt808-ingest-app/pom.xml | 2 +- modules/apps/vehicle-analytics-app/pom.xml | 2 +- modules/apps/vehicle-history-app/pom.xml | 2 +- .../historyapp/MavenModuleProfileTest.java | 21 +++++++++ modules/apps/xinda-push-app/pom.xml | 2 +- modules/apps/yutong-mqtt-app/pom.xml | 2 +- .../services/event-history-service/pom.xml | 2 +- modules/services/vehicle-stat-service/pom.xml | 2 +- .../services/vehicle-state-service/pom.xml | 2 +- modules/sinks/{sink-mq => sink-kafka}/pom.xml | 4 +- .../ingest/sink/kafka/EnvelopeMapper.java | 0 .../kafka/KafkaEnvelopeConsumerFactory.java | 0 .../kafka/KafkaEnvelopeConsumerRunner.java | 0 .../kafka/KafkaEnvelopeConsumerWorker.java | 0 .../kafka/KafkaEnvelopeDeadLetterSink.java | 0 .../ingest/sink/kafka/KafkaEventSink.java | 0 .../kafka/KafkaSinkAutoConfiguration.java | 0 .../KafkaSinkConsumerAutoConfiguration.java | 0 .../sink/kafka/KafkaSinkProperties.java | 0 .../ingest/sink/kafka/TopicRouter.java | 0 .../src/main/proto/vehicle_envelope.proto | 0 ...ot.autoconfigure.AutoConfiguration.imports | 0 .../EnvelopeMapperTelemetrySnapshotTest.java | 0 .../KafkaEnvelopeConsumerFactoryTest.java | 0 .../KafkaEnvelopeConsumerRunnerTest.java | 0 .../KafkaEnvelopeConsumerWorkerTest.java | 0 .../KafkaEnvelopeDeadLetterSinkTest.java | 0 .../ingest/sink/kafka/KafkaEventSinkTest.java | 0 ...afkaSinkConsumerAutoConfigurationTest.java | 0 .../sink/kafka/KafkaSinkPropertiesTest.java | 0 .../ingest/sink/kafka/TopicRouterTest.java | 0 ...VehicleEnvelopeProtoCompatibilityTest.java | 0 modules/sinks/tdengine-history-store/pom.xml | 2 +- pom.xml | 4 +- 43 files changed, 87 insertions(+), 66 deletions(-) rename modules/sinks/{sink-mq => sink-kafka}/pom.xml (98%) rename modules/sinks/{sink-mq => sink-kafka}/src/main/java/com/lingniu/ingest/sink/kafka/EnvelopeMapper.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerFactory.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerRunner.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerWorker.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeDeadLetterSink.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEventSink.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkAutoConfiguration.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfiguration.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkProperties.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/main/java/com/lingniu/ingest/sink/kafka/TopicRouter.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/main/proto/vehicle_envelope.proto (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/test/java/com/lingniu/ingest/sink/kafka/EnvelopeMapperTelemetrySnapshotTest.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerFactoryTest.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerRunnerTest.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerWorkerTest.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeDeadLetterSinkTest.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEventSinkTest.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfigurationTest.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkPropertiesTest.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/test/java/com/lingniu/ingest/sink/kafka/TopicRouterTest.java (100%) rename modules/sinks/{sink-mq => sink-kafka}/src/test/java/com/lingniu/ingest/sink/kafka/VehicleEnvelopeProtoCompatibilityTest.java (100%) diff --git a/CHANGELOG.md b/CHANGELOG.md index c24a658b..7d5b33be 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -56,7 +56,7 @@ `protocol-jsatl12`)。 - **MQTT 入站**与**信达平台 push 接入**(`inbound-mqtt` / `inbound-xinda-push`)。 - **Disruptor 事件总线**(`ingest-core`):高吞吐解耦解码与下游 sink。 -- **Sink**:本地归档(`sink-archive`)、Kafka(`sink-mq`,含 protobuf envelope)。 +- **Sink**:本地归档(`sink-archive`)、Kafka(`sink-kafka`,含 protobuf envelope)。 - **会话状态**(`session-core`)。 - **可观测性**:Micrometer/Actuator 装配(`observability`)。 - **终端控制命令网关**(`command-gateway`)。 diff --git a/README.md b/README.md index 5a476fd9..16925080 100644 --- a/README.md +++ b/README.md @@ -48,7 +48,7 @@ lingniu-vehicle-ingest/ │ │ ├── inbound-mqtt/ MQTT 接入(endpoint 生命周期 + profile 注册扩展 + 统一身份映射 + PEM 双向 TLS + 未知 profile/解析失败/profile异常/连接订阅失败兜底 + 统一 UNKNOWN 身份 metadata) │ │ └── inbound-xinda-push/ 信达 Push 接入(废弃兼容模块,仅 -Plegacy-xinda 显式构建) │ ├── sinks/ -│ │ ├── sink-mq/ Kafka producer + Protobuf Envelope +│ │ ├── sink-kafka/ Kafka producer + Protobuf Envelope │ │ ├── sink-archive/ 原始报文冷存 │ │ └── event-file-store/ 可选兼容索引(旧 Parquet/DuckDB 路径,仅 -Poptional-event-file-store 显式构建) │ ├── services/ diff --git a/docs/module-data-flow.html b/docs/module-data-flow.html index b59acfc2..e7d29739 100644 --- a/docs/module-data-flow.html +++ b/docs/module-data-flow.html @@ -317,7 +317,7 @@ flowchart TB end subgraph sink["输出与明细存储 modules/sinks"] - mq["sink-mq
VehicleEvent 到 Protobuf Envelope 到 Kafka"] + mq["sink-kafka
VehicleEvent 到 Protobuf Envelope 到 Kafka"] archive["sink-archive
RawArchive 到本地/S3/OSS 冷存"] tdengineStore["tdengine-history-store
TDengine raw_frames + locations"] fileStore["event-file-store
Parquet/DuckDB 兼容可选"] @@ -416,7 +416,7 @@ sequenceDiagram participant Dispatcher as Dispatcher participant Handler as Protocol Handler
Mapper participant EventBus as DisruptorEventBus - participant MQ as sink-mq
Kafka Envelope + participant MQ as sink-kafka
Kafka Envelope participant Archive as sink-archive
ArchiveStore participant History as vehicle-history-app participant TDengine as tdengine-history-store
raw_frames + locations @@ -578,7 +578,7 @@ flowchart LR 车企平台接入profile 扩展身份映射双向 TLS错误隔离 - modules/sinks/sink-mq + modules/sinks/sink-kafka Kafka Sink。把 VehicleEvent 转成 Protobuf Envelope,按 VIN 分区投递;业务事件携带全字段 TelemetrySnapshot 和 rawArchiveUri,显式 RawArchive Envelope 会填充 archive:// 逻辑 URI 与 size,默认 Kafka Sink 仍不发送原始 bytes 本体;同时提供 KafkaEnvelopeConsumerFactory、KafkaEnvelopeConsumerRunner、KafkaEnvelopeConsumerWorker 和 KafkaEnvelopeDeadLetterSink,统一服务侧 Kafka 消费启动、独立 group 绑定、死信发布和 topic 到 EnvelopeConsumerProcessor 的分发。 VehicleEvent。 Kafka Protobuf 消息、消费侧 DLQ 记录。 diff --git a/docs/operations/gb32960-service-split-runbook.md b/docs/operations/gb32960-service-split-runbook.md index ff601645..c310a468 100644 --- a/docs/operations/gb32960-service-split-runbook.md +++ b/docs/operations/gb32960-service-split-runbook.md @@ -251,7 +251,7 @@ Observed on 2026-06-23 in worktree `.worktrees/gb32960-service-split`: Observed on 2026-06-23 in worktree `.worktrees/gb32960-service-split`: -- Targeted module tests: `mvn -pl :sink-mq,:ingest-core,:protocol-gb32960,:event-history-service,:vehicle-state-service,:vehicle-stat-service test` ended with `BUILD SUCCESS`; Maven reported 55 protocol/sink/core tests, 30 event-history tests, 14 vehicle-state tests, and 19 vehicle-stat tests with 0 failures/errors/skips. +- Targeted module tests: `mvn -pl :sink-kafka,:ingest-core,:protocol-gb32960,:event-history-service,:vehicle-state-service,:vehicle-stat-service test` ended with `BUILD SUCCESS`; Maven reported 55 protocol/sink/core tests, 30 event-history tests, 14 vehicle-state tests, and 19 vehicle-stat tests with 0 failures/errors/skips. - Split app tests: `mvn -pl :gb32960-ingest-app,:vehicle-history-app,:vehicle-analytics-app test` ended with `BUILD SUCCESS`; Maven reported 6 app composition/default tests with 0 failures/errors/skips. - Package verification: `mvn -pl :gb32960-ingest-app,:vehicle-history-app,:vehicle-analytics-app -am package -Dmaven.test.skip=true` ended with `BUILD SUCCESS` and repackaged all three app jars. - Split app startup: not run in this verification pass because local Kafka was still absent. diff --git a/docs/superpowers/plans/2026-04-20-gb32960-body-parser-block-isolation.md b/docs/superpowers/plans/2026-04-20-gb32960-body-parser-block-isolation.md index 9149676b..ad1c2ff3 100644 --- a/docs/superpowers/plans/2026-04-20-gb32960-body-parser-block-isolation.md +++ b/docs/superpowers/plans/2026-04-20-gb32960-body-parser-block-isolation.md @@ -543,7 +543,7 @@ EOF Run: `mvn test -q` Expected: BUILD SUCCESS。重点关注: - `protocol-gb32960` 模块全绿(Golden、FullBlocks、Isolation、GuangdongFcEndToEnd、Decoder、Mapper、所有 parser 单测、profile) -- 其它模块(`ingest-core`、`sink-mq` 等)不受影响(无改动应自动通过) +- 其它模块(`ingest-core`、`sink-kafka` 等)不受影响(无改动应自动通过) 若有红,停下来诊断。 diff --git a/docs/superpowers/plans/2026-06-23-gb32960-service-split.md b/docs/superpowers/plans/2026-06-23-gb32960-service-split.md index 6b2e2873..84f02c7a 100644 --- a/docs/superpowers/plans/2026-06-23-gb32960-service-split.md +++ b/docs/superpowers/plans/2026-06-23-gb32960-service-split.md @@ -48,11 +48,11 @@ Create app modules: Modify shared build and Kafka contract files: - `pom.xml` -- `modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaSinkProperties.java` -- `modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/TopicRouter.java` -- `modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEventSink.java` -- `modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/TopicRouterTest.java` -- `modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEventSinkTest.java` +- `modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkProperties.java` +- `modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/TopicRouter.java` +- `modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEventSink.java` +- `modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/TopicRouterTest.java` +- `modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEventSinkTest.java` Modify ACK/Kafka durability boundary files: @@ -157,7 +157,7 @@ Create `modules/apps/gb32960-ingest-app/pom.xml`: com.lingniu.ingest - sink-mq + sink-kafka org.springframework.boot @@ -222,7 +222,7 @@ Create `modules/apps/vehicle-history-app/pom.xml`: com.lingniu.ingest - sink-mq + sink-kafka com.lingniu.ingest @@ -303,7 +303,7 @@ Create `modules/apps/vehicle-analytics-app/pom.xml`: com.lingniu.ingest - sink-mq + sink-kafka com.lingniu.ingest @@ -933,15 +933,15 @@ git commit -m "test: cover split service composition" ### Task 4: Version Kafka Topics for GB32960 Raw/Event/DLQ **Files:** -- Modify: `modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaSinkProperties.java` -- Modify: `modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/TopicRouter.java` -- Modify: `modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEventSink.java` -- Create or modify: `modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/TopicRouterTest.java` -- Create or modify: `modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEventSinkTest.java` +- Modify: `modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkProperties.java` +- Modify: `modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/TopicRouter.java` +- Modify: `modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEventSink.java` +- Create or modify: `modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/TopicRouterTest.java` +- Create or modify: `modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEventSinkTest.java` - [ ] **Step 1: Write topic routing tests** -Create `modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/TopicRouterTest.java` if it does not exist. Add tests proving GB32960 raw and normalized events can route to versioned topics. Use the existing `VehicleEvent` constructors from `EventFileStoreSinkTest` and `EnvelopeMapperTelemetrySnapshotTest` as examples. +Create `modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/TopicRouterTest.java` if it does not exist. Add tests proving GB32960 raw and normalized events can route to versioned topics. Use the existing `VehicleEvent` constructors from `EventFileStoreSinkTest` and `EnvelopeMapperTelemetrySnapshotTest` as examples. Expected test names: @@ -996,7 +996,7 @@ Adjust constructor arguments to match the current sealed `VehicleEvent` definiti Run: ```bash -mvn -pl :sink-mq -Dtest=TopicRouterTest test +mvn -pl :sink-kafka -Dtest=TopicRouterTest test ``` Expected before implementation: compile failure if constructors are wrong or assertion failure if defaults are old. Fix constructors first; keep routing assertions. @@ -1028,12 +1028,12 @@ public boolean accepts(VehicleEvent event) { Update the JavaDoc above it to say Kafka carries both normalized events and raw archive records in split-service mode. -- [ ] **Step 5: Run sink-mq tests** +- [ ] **Step 5: Run sink-kafka tests** Run: ```bash -mvn -pl :sink-mq test +mvn -pl :sink-kafka test ``` Expected: PASS. @@ -1041,11 +1041,11 @@ Expected: PASS. - [ ] **Step 6: Commit** ```bash -git add modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaSinkProperties.java \ - modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/TopicRouter.java \ - modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEventSink.java \ - modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/TopicRouterTest.java \ - modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEventSinkTest.java +git add modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkProperties.java \ + modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/TopicRouter.java \ + modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEventSink.java \ + modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/TopicRouterTest.java \ + modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEventSinkTest.java git commit -m "feat: route gb32960 kafka topics" ``` @@ -1567,7 +1567,7 @@ git commit -m "docs: add gb32960 split service runbook" Run: ```bash -mvn -pl :sink-mq,:ingest-core,:protocol-gb32960,:event-history-service,:vehicle-state-service,:vehicle-stat-service test +mvn -pl :sink-kafka,:ingest-core,:protocol-gb32960,:event-history-service,:vehicle-state-service,:vehicle-stat-service test ``` Expected: BUILD SUCCESS. diff --git a/docs/superpowers/plans/2026-06-29-vehicle-ingest-redesign-phase1.md b/docs/superpowers/plans/2026-06-29-vehicle-ingest-redesign-phase1.md index 36e82acf..b34cb282 100644 --- a/docs/superpowers/plans/2026-06-29-vehicle-ingest-redesign-phase1.md +++ b/docs/superpowers/plans/2026-06-29-vehicle-ingest-redesign-phase1.md @@ -16,7 +16,7 @@ This plan implements only the first slice of the redesign: - `modules/core/ingest-facts`: protocol-neutral fact records and helpers. - `modules/sinks/raw-archive-store`: raw archive write/read contract and local URI-safe implementation skeleton. -- `modules/sinks/sink-mq`: protobuf messages for raw frame facts and decoded facts. +- `modules/sinks/sink-kafka`: protobuf messages for raw frame facts and decoded facts. - Parent Maven wiring and focused tests. It intentionally does not yet modify: @@ -53,9 +53,9 @@ Create: Modify: - `pom.xml`: add modules and dependency management entries. -- `modules/sinks/sink-mq/pom.xml`: add dependency on `ingest-facts`. -- `modules/sinks/sink-mq/src/main/proto/vehicle_envelope.proto`: add `RawFrameFactPayload`, `DecodedFactPayload`, and related enums/messages. -- `modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/VehicleEnvelopeProtoCompatibilityTest.java`: verify new proto messages can round-trip. +- `modules/sinks/sink-kafka/pom.xml`: add dependency on `ingest-facts`. +- `modules/sinks/sink-kafka/src/main/proto/vehicle_envelope.proto`: add `RawFrameFactPayload`, `DecodedFactPayload`, and related enums/messages. +- `modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/VehicleEnvelopeProtoCompatibilityTest.java`: verify new proto messages can round-trip. --- @@ -1260,9 +1260,9 @@ git commit -m "feat: add raw archive store contract" **Files:** -- Modify: `modules/sinks/sink-mq/pom.xml` -- Modify: `modules/sinks/sink-mq/src/main/proto/vehicle_envelope.proto` -- Create: `modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/VehicleEnvelopeProtoCompatibilityTest.java` +- Modify: `modules/sinks/sink-kafka/pom.xml` +- Modify: `modules/sinks/sink-kafka/src/main/proto/vehicle_envelope.proto` +- Create: `modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/VehicleEnvelopeProtoCompatibilityTest.java` - [ ] **Step 1: Write failing proto compatibility test** @@ -1341,14 +1341,14 @@ class VehicleEnvelopeProtoCompatibilityTest { Run: ```bash -mvn -pl modules/sinks/sink-mq -Dtest=VehicleEnvelopeProtoCompatibilityTest test +mvn -pl modules/sinks/sink-kafka -Dtest=VehicleEnvelopeProtoCompatibilityTest test ``` Expected: FAIL because proto generated classes do not exist. -- [ ] **Step 3: Add `ingest-facts` dependency to sink-mq** +- [ ] **Step 3: Add `ingest-facts` dependency to sink-kafka** -Add this dependency to `modules/sinks/sink-mq/pom.xml` after `ingest-api`: +Add this dependency to `modules/sinks/sink-kafka/pom.xml` after `ingest-api`: ```xml @@ -1412,17 +1412,17 @@ message DecodedFactPayload { Run: ```bash -mvn -pl modules/sinks/sink-mq -Dtest=VehicleEnvelopeProtoCompatibilityTest test +mvn -pl modules/sinks/sink-kafka -Dtest=VehicleEnvelopeProtoCompatibilityTest test ``` Expected: PASS. -- [ ] **Step 6: Run existing sink-mq tests** +- [ ] **Step 6: Run existing sink-kafka tests** Run: ```bash -mvn -pl modules/sinks/sink-mq test +mvn -pl modules/sinks/sink-kafka test ``` Expected: PASS. Existing envelope tests must continue to pass because new proto fields are additive. @@ -1430,9 +1430,9 @@ Expected: PASS. Existing envelope tests must continue to pass because new proto - [ ] **Step 7: Commit** ```bash -git add modules/sinks/sink-mq/pom.xml \ - modules/sinks/sink-mq/src/main/proto/vehicle_envelope.proto \ - modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/VehicleEnvelopeProtoCompatibilityTest.java +git add modules/sinks/sink-kafka/pom.xml \ + modules/sinks/sink-kafka/src/main/proto/vehicle_envelope.proto \ + modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/VehicleEnvelopeProtoCompatibilityTest.java git commit -m "feat: add fact payloads to kafka envelope" ``` @@ -1449,7 +1449,7 @@ git commit -m "feat: add fact payloads to kafka envelope" Run: ```bash -mvn -pl modules/core/ingest-facts,modules/sinks/raw-archive-store,modules/sinks/sink-mq -am test +mvn -pl modules/core/ingest-facts,modules/sinks/raw-archive-store,modules/sinks/sink-kafka -am test ``` Expected: PASS. @@ -1459,7 +1459,7 @@ Expected: PASS. Run: ```bash -mvn -pl modules/sinks/sink-mq -am -DskipTests package +mvn -pl modules/sinks/sink-kafka -am -DskipTests package ``` Expected: PASS. Generated protobuf Java sources include `RawFrameFactPayload`, `DecodedFactPayload`, and `ParseStatusProto`. @@ -1473,7 +1473,7 @@ rg -n "RawFrameFact|DecodedFact|RawArchiveWriter|vehicle.raw-frame.v1|vehicle.de modules/apps modules/protocols modules/services modules/core modules/sinks ``` -Expected: Matches only in the new modules, sink-mq proto/test, and planned references. No existing handler behavior should be changed in Phase 1. +Expected: Matches only in the new modules, sink-kafka proto/test, and planned references. No existing handler behavior should be changed in Phase 1. - [ ] **Step 4: Commit final verification note if any docs changed** diff --git a/docs/superpowers/specs/2026-06-23-gb32960-service-split-design.md b/docs/superpowers/specs/2026-06-23-gb32960-service-split-design.md index d6d386fb..a3f84303 100644 --- a/docs/superpowers/specs/2026-06-23-gb32960-service-split-design.md +++ b/docs/superpowers/specs/2026-06-23-gb32960-service-split-design.md @@ -60,7 +60,7 @@ Initial module dependencies: - `ingest-codec-common` - `session-core` - `protocol-gb32960` -- `sink-mq` +- `sink-kafka` - `observability` ### `vehicle-history-app` @@ -86,7 +86,7 @@ Non-responsibilities: Initial module dependencies: - `ingest-api` -- `sink-mq` +- `sink-kafka` - `sink-archive` - `event-file-store` - `event-history-service` @@ -116,7 +116,7 @@ Non-responsibilities: Initial module dependencies: - `ingest-api` -- `sink-mq` +- `sink-kafka` - `vehicle-state-service` - `vehicle-stat-service` - `observability` diff --git a/modules/apps/gb32960-ingest-app/pom.xml b/modules/apps/gb32960-ingest-app/pom.xml index dda16f3e..d9652887 100644 --- a/modules/apps/gb32960-ingest-app/pom.xml +++ b/modules/apps/gb32960-ingest-app/pom.xml @@ -35,7 +35,7 @@ com.lingniu.ingest - sink-mq + sink-kafka com.lingniu.ingest diff --git a/modules/apps/jt808-ingest-app/pom.xml b/modules/apps/jt808-ingest-app/pom.xml index 30cdfbd6..240aa925 100644 --- a/modules/apps/jt808-ingest-app/pom.xml +++ b/modules/apps/jt808-ingest-app/pom.xml @@ -35,7 +35,7 @@ com.lingniu.ingest - sink-mq + sink-kafka com.lingniu.ingest diff --git a/modules/apps/vehicle-analytics-app/pom.xml b/modules/apps/vehicle-analytics-app/pom.xml index 98de9d4a..559cf965 100644 --- a/modules/apps/vehicle-analytics-app/pom.xml +++ b/modules/apps/vehicle-analytics-app/pom.xml @@ -23,7 +23,7 @@ com.lingniu.ingest - sink-mq + sink-kafka com.lingniu.ingest diff --git a/modules/apps/vehicle-history-app/pom.xml b/modules/apps/vehicle-history-app/pom.xml index 07c1ea71..6df80520 100644 --- a/modules/apps/vehicle-history-app/pom.xml +++ b/modules/apps/vehicle-history-app/pom.xml @@ -23,7 +23,7 @@ com.lingniu.ingest - sink-mq + sink-kafka com.lingniu.ingest diff --git a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/MavenModuleProfileTest.java b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/MavenModuleProfileTest.java index d4fa8568..e17f0594 100644 --- a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/MavenModuleProfileTest.java +++ b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/MavenModuleProfileTest.java @@ -79,6 +79,27 @@ class MavenModuleProfileTest { .isFalse(); } + @Test + void kafkaSinkModuleUsesKafkaNamingInsteadOfMqNaming() throws Exception { + Document pom = rootPom(); + Document kafkaSinkPom = modulePom("modules/sinks/sink-kafka/pom.xml"); + String readme = Files.readString(repositoryRoot().resolve("README.md")); + String legacyModule = "modules/sinks/sink-" + "mq"; + String legacyArtifact = "sink-" + "mq"; + + assertThat(defaultModules(pom)) + .contains("modules/sinks/sink-kafka") + .doesNotContain(legacyModule); + assertThat(rootDependencyManagementArtifacts(pom)) + .contains("sink-kafka") + .doesNotContain(legacyArtifact); + assertThat(firstDirectChild(kafkaSinkPom.getDocumentElement(), "artifactId").getTextContent()) + .isEqualTo("sink-kafka"); + assertThat(readme) + .contains("sink-kafka/") + .doesNotContain(legacyArtifact); + } + @Test void xindaIsLegacyOnlyAndOutsideDefaultProductionReactor() throws Exception { Document pom = rootPom(); diff --git a/modules/apps/xinda-push-app/pom.xml b/modules/apps/xinda-push-app/pom.xml index 487c7930..052910c6 100644 --- a/modules/apps/xinda-push-app/pom.xml +++ b/modules/apps/xinda-push-app/pom.xml @@ -31,7 +31,7 @@ com.lingniu.ingest - sink-mq + sink-kafka com.lingniu.ingest diff --git a/modules/apps/yutong-mqtt-app/pom.xml b/modules/apps/yutong-mqtt-app/pom.xml index 70c8ca1b..6e515674 100644 --- a/modules/apps/yutong-mqtt-app/pom.xml +++ b/modules/apps/yutong-mqtt-app/pom.xml @@ -31,7 +31,7 @@ com.lingniu.ingest - sink-mq + sink-kafka com.lingniu.ingest diff --git a/modules/services/event-history-service/pom.xml b/modules/services/event-history-service/pom.xml index 19cfc4e4..e9065f59 100644 --- a/modules/services/event-history-service/pom.xml +++ b/modules/services/event-history-service/pom.xml @@ -22,7 +22,7 @@ com.lingniu.ingest - sink-mq + sink-kafka com.lingniu.ingest diff --git a/modules/services/vehicle-stat-service/pom.xml b/modules/services/vehicle-stat-service/pom.xml index 28c1268f..4855e022 100644 --- a/modules/services/vehicle-stat-service/pom.xml +++ b/modules/services/vehicle-stat-service/pom.xml @@ -14,7 +14,7 @@ com.lingniu.ingest - sink-mq + sink-kafka org.springframework.boot diff --git a/modules/services/vehicle-state-service/pom.xml b/modules/services/vehicle-state-service/pom.xml index 0413b743..360dcd44 100644 --- a/modules/services/vehicle-state-service/pom.xml +++ b/modules/services/vehicle-state-service/pom.xml @@ -14,7 +14,7 @@ com.lingniu.ingest - sink-mq + sink-kafka org.springframework.boot diff --git a/modules/sinks/sink-mq/pom.xml b/modules/sinks/sink-kafka/pom.xml similarity index 98% rename from modules/sinks/sink-mq/pom.xml rename to modules/sinks/sink-kafka/pom.xml index 03a9d565..ca0a1329 100644 --- a/modules/sinks/sink-mq/pom.xml +++ b/modules/sinks/sink-kafka/pom.xml @@ -7,8 +7,8 @@ 0.1.0-SNAPSHOT ../../../pom.xml - sink-mq - sink-mq + sink-kafka + sink-kafka Kafka producer + Protobuf Envelope + 重试/熔断/DLQ。 diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/EnvelopeMapper.java b/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/EnvelopeMapper.java similarity index 100% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/EnvelopeMapper.java rename to modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/EnvelopeMapper.java diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerFactory.java b/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerFactory.java similarity index 100% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerFactory.java rename to modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerFactory.java diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerRunner.java b/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerRunner.java similarity index 100% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerRunner.java rename to modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerRunner.java diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerWorker.java b/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerWorker.java similarity index 100% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerWorker.java rename to modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerWorker.java diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeDeadLetterSink.java b/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeDeadLetterSink.java similarity index 100% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeDeadLetterSink.java rename to modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeDeadLetterSink.java diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEventSink.java b/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEventSink.java similarity index 100% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEventSink.java rename to modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEventSink.java diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkAutoConfiguration.java b/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkAutoConfiguration.java similarity index 100% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkAutoConfiguration.java rename to modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkAutoConfiguration.java diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfiguration.java b/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfiguration.java similarity index 100% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfiguration.java rename to modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfiguration.java diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkProperties.java b/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkProperties.java similarity index 100% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkProperties.java rename to modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkProperties.java diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/TopicRouter.java b/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/TopicRouter.java similarity index 100% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/TopicRouter.java rename to modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/TopicRouter.java diff --git a/modules/sinks/sink-mq/src/main/proto/vehicle_envelope.proto b/modules/sinks/sink-kafka/src/main/proto/vehicle_envelope.proto similarity index 100% rename from modules/sinks/sink-mq/src/main/proto/vehicle_envelope.proto rename to modules/sinks/sink-kafka/src/main/proto/vehicle_envelope.proto diff --git a/modules/sinks/sink-mq/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports b/modules/sinks/sink-kafka/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports similarity index 100% rename from modules/sinks/sink-mq/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports rename to modules/sinks/sink-kafka/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/EnvelopeMapperTelemetrySnapshotTest.java b/modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/EnvelopeMapperTelemetrySnapshotTest.java similarity index 100% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/EnvelopeMapperTelemetrySnapshotTest.java rename to modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/EnvelopeMapperTelemetrySnapshotTest.java diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerFactoryTest.java b/modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerFactoryTest.java similarity index 100% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerFactoryTest.java rename to modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerFactoryTest.java diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerRunnerTest.java b/modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerRunnerTest.java similarity index 100% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerRunnerTest.java rename to modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerRunnerTest.java diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerWorkerTest.java b/modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerWorkerTest.java similarity index 100% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerWorkerTest.java rename to modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerWorkerTest.java diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeDeadLetterSinkTest.java b/modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeDeadLetterSinkTest.java similarity index 100% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeDeadLetterSinkTest.java rename to modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeDeadLetterSinkTest.java diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEventSinkTest.java b/modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEventSinkTest.java similarity index 100% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEventSinkTest.java rename to modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEventSinkTest.java diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfigurationTest.java b/modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfigurationTest.java similarity index 100% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfigurationTest.java rename to modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfigurationTest.java diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkPropertiesTest.java b/modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkPropertiesTest.java similarity index 100% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkPropertiesTest.java rename to modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkPropertiesTest.java diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/TopicRouterTest.java b/modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/TopicRouterTest.java similarity index 100% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/TopicRouterTest.java rename to modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/TopicRouterTest.java diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/VehicleEnvelopeProtoCompatibilityTest.java b/modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/VehicleEnvelopeProtoCompatibilityTest.java similarity index 100% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/VehicleEnvelopeProtoCompatibilityTest.java rename to modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/VehicleEnvelopeProtoCompatibilityTest.java diff --git a/modules/sinks/tdengine-history-store/pom.xml b/modules/sinks/tdengine-history-store/pom.xml index f5798f52..be3ea689 100644 --- a/modules/sinks/tdengine-history-store/pom.xml +++ b/modules/sinks/tdengine-history-store/pom.xml @@ -14,7 +14,7 @@ com.lingniu.ingest - sink-mq + sink-kafka org.springframework.boot diff --git a/pom.xml b/pom.xml index 471099e9..4772b126 100644 --- a/pom.xml +++ b/pom.xml @@ -24,7 +24,7 @@ modules/core/session-core modules/core/vehicle-identity modules/core/observability - modules/sinks/sink-mq + modules/sinks/sink-kafka modules/sinks/sink-archive modules/sinks/tdengine-history-store modules/services/event-history-service @@ -279,7 +279,7 @@ com.lingniu.ingest - sink-mq + sink-kafka ${project.version}