From 849e8f8a8b73f962bbc437c647d9d000e2543ed2 Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 10:50:50 +0800 Subject: [PATCH] refactor: expose kafka sink configuration --- .../plans/2026-06-23-gb32960-service-split.md | 53 ++++++++--------- ...26-06-29-vehicle-ingest-redesign-phase1.md | 10 ++-- ...2026-06-23-gb32960-service-split-design.md | 2 +- docs/target-architecture.md | 2 +- .../src/main/resources/application.yml | 3 +- .../Gb32960IngestAppCompositionTest.java | 15 +++-- .../Gb32960IngestAppDefaultsTest.java | 10 +++- .../src/main/resources/application.yml | 3 +- .../Jt808IngestAppCompositionTest.java | 15 +++-- .../jt808app/Jt808IngestAppDefaultsTest.java | 10 ++-- .../src/main/resources/application.yml | 3 +- .../VehicleAnalyticsAppCompositionTest.java | 15 +++-- .../VehicleAnalyticsAppDefaultsTest.java | 6 +- ...icleHistoryKafkaConsumerConfiguration.java | 12 ++-- .../src/main/resources/application.yml | 3 +- .../VehicleHistoryAppCompositionTest.java | 58 +++++++++---------- .../VehicleHistoryAppDefaultsTest.java | 44 +++++++------- .../src/main/resources/application.yml | 16 ++--- .../XindaPushAppCompositionTest.java | 17 +++--- .../XindaPushAppDefaultsTest.java | 15 ++--- .../src/main/resources/application.yml | 3 +- .../YutongMqttAppCompositionTest.java | 15 +++-- .../YutongMqttAppDefaultsTest.java | 10 ++-- .../EventHistoryEnvelopeIngestor.java | 2 +- .../TelemetryEnvelopeRecordMapper.java | 4 +- .../config/EventHistoryAutoConfiguration.java | 8 +-- .../EventHistoryEnvelopeIngestorTest.java | 14 ++--- .../TelemetryEnvelopeRecordMapperTest.java | 8 +-- .../EventHistoryAutoConfigurationTest.java | 2 +- .../config/EventHistoryPomBoundaryTest.java | 2 +- .../VehicleStatEnvelopeIngestor.java | 2 +- .../config/VehicleStatAutoConfiguration.java | 2 +- .../jt808/Jt808LocationPointExtractor.java | 4 +- .../jt808/Jt808MileageStreamProcessor.java | 2 +- .../VehicleStatEnvelopeIngestorTest.java | 8 +-- .../Jt808LocationPointExtractorTest.java | 6 +- .../Jt808MileageStreamProcessorTest.java | 6 +- .../VehicleStateEnvelopeIngestor.java | 2 +- .../vehiclestate/VehicleStateUpdater.java | 4 +- .../config/VehicleStateAutoConfiguration.java | 2 +- .../VehicleStateEnvelopeIngestorTest.java | 8 +-- .../vehiclestate/VehicleStateUpdaterTest.java | 6 +- .../sink/{mq => kafka}/EnvelopeMapper.java | 40 ++++++------- .../KafkaEnvelopeConsumerFactory.java | 30 +++++----- .../KafkaEnvelopeConsumerRunner.java | 2 +- .../KafkaEnvelopeConsumerWorker.java | 2 +- .../KafkaEnvelopeDeadLetterSink.java | 2 +- .../sink/{mq => kafka}/KafkaEventSink.java | 2 +- .../KafkaSinkAutoConfiguration.java} | 38 +++++------- .../KafkaSinkConsumerAutoConfiguration.java} | 16 ++--- .../KafkaSinkProperties.java} | 21 ++----- .../sink/{mq => kafka}/TopicRouter.java | 6 +- .../src/main/proto/vehicle_envelope.proto | 4 +- ...ot.autoconfigure.AutoConfiguration.imports | 4 +- .../EnvelopeMapperTelemetrySnapshotTest.java | 4 +- .../KafkaEnvelopeConsumerFactoryTest.java | 6 +- .../KafkaEnvelopeConsumerRunnerTest.java | 2 +- .../KafkaEnvelopeConsumerWorkerTest.java | 2 +- .../KafkaEnvelopeDeadLetterSinkTest.java | 2 +- .../{mq => kafka}/KafkaEventSinkTest.java | 6 +- ...fkaSinkConsumerAutoConfigurationTest.java} | 19 +++--- .../sink/kafka/KafkaSinkPropertiesTest.java | 26 +++++++++ .../sink/{mq => kafka}/TopicRouterTest.java | 8 +-- ...VehicleEnvelopeProtoCompatibilityTest.java | 10 ++-- .../ingest/sink/mq/SinkMqPropertiesTest.java | 29 ---------- .../tdenginehistory/TdengineEnvelopeRows.java | 8 +-- .../TdengineEnvelopeRowsTest.java | 14 ++--- 67 files changed, 350 insertions(+), 385 deletions(-) rename modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/{mq => kafka}/EnvelopeMapper.java (86%) rename modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/{mq => kafka}/KafkaEnvelopeConsumerFactory.java (83%) rename modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/{mq => kafka}/KafkaEnvelopeConsumerRunner.java (99%) rename modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/{mq => kafka}/KafkaEnvelopeConsumerWorker.java (98%) rename modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/{mq => kafka}/KafkaEnvelopeDeadLetterSink.java (98%) rename modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/{mq => kafka}/KafkaEventSink.java (99%) rename modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/{mq/SinkMqAutoConfiguration.java => kafka/KafkaSinkAutoConfiguration.java} (64%) rename modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/{mq/SinkMqConsumerAutoConfiguration.java => kafka/KafkaSinkConsumerAutoConfiguration.java} (78%) rename modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/{mq/SinkMqProperties.java => kafka/KafkaSinkProperties.java} (90%) rename modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/{mq => kafka}/TopicRouter.java (84%) rename modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/{mq => kafka}/EnvelopeMapperTelemetrySnapshotTest.java (97%) rename modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/{mq => kafka}/KafkaEnvelopeConsumerFactoryTest.java (95%) rename modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/{mq => kafka}/KafkaEnvelopeConsumerRunnerTest.java (97%) rename modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/{mq => kafka}/KafkaEnvelopeConsumerWorkerTest.java (99%) rename modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/{mq => kafka}/KafkaEnvelopeDeadLetterSinkTest.java (98%) rename modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/{mq => kafka}/KafkaEventSinkTest.java (94%) rename modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/{mq/SinkMqConsumerAutoConfigurationTest.java => kafka/KafkaSinkConsumerAutoConfigurationTest.java} (70%) create mode 100644 modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkPropertiesTest.java rename modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/{mq => kafka}/TopicRouterTest.java (94%) rename modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/{mq => kafka}/VehicleEnvelopeProtoCompatibilityTest.java (90%) delete mode 100644 modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/SinkMqPropertiesTest.java 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 0cbd3ae3..6b2e2873 100644 --- a/docs/superpowers/plans/2026-06-23-gb32960-service-split.md +++ b/docs/superpowers/plans/2026-06-23-gb32960-service-split.md @@ -18,7 +18,7 @@ - Existing handler currently writes report ACK before `dispatcher.dispatch(rf)`. - Existing `DisruptorEventBus` publishes to sinks asynchronously and does not wait for sink completion. - Existing `KafkaEventSink.accepts()` rejects `VehicleEvent.RawArchive`. -- Existing versioned topic names are not yet present; current topic names live under `lingniu.ingest.sink.mq.topics`. +- Existing versioned topic names are not yet present; current topic names live under `lingniu.ingest.sink.kafka.topics`. - Existing history consumer logic: `modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestor.java`. - Existing analytics consumers: `VehicleStateEnvelopeIngestor` and `VehicleStatEnvelopeIngestor`. - There is an unrelated working-tree change in `modules/apps/bootstrap-all/src/main/resources/application.yml` for GB32960 diagnostics capacity. Do not revert it. Do not include it in service-split commits unless the user explicitly asks. @@ -48,7 +48,7 @@ Create app modules: Modify shared build and Kafka contract files: - `pom.xml` -- `modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqProperties.java` +- `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` @@ -466,16 +466,19 @@ lingniu: rate-limit: per-vin-qps: 50 session: - store: ${SESSION_STORE:memory} + store: ${SESSION_STORE:redis} ttl: ${SESSION_TTL:30m} identity: - store: ${VEHICLE_IDENTITY_STORE:file} - file: - path: ${VEHICLE_IDENTITY_FILE:./data/vehicle-identity.jsonl} + store: ${VEHICLE_IDENTITY_STORE:mysql} + mysql: + jdbc-url: ${VEHICLE_IDENTITY_MYSQL_JDBC_URL:jdbc:mysql://127.0.0.1:3306/lingniu_vehicle?useUnicode=true&characterEncoding=utf8&useSSL=false&serverTimezone=Asia/Shanghai} + username: ${VEHICLE_IDENTITY_MYSQL_USERNAME:root} + password: ${VEHICLE_IDENTITY_MYSQL_PASSWORD:} + table-name: ${VEHICLE_IDENTITY_MYSQL_TABLE:vehicle_identity_bindings} + initialize-schema: ${VEHICLE_IDENTITY_MYSQL_INITIALIZE_SCHEMA:true} sink: - mq: + kafka: enabled: ${KAFKA_ENABLED:true} - type: kafka bootstrap-servers: ${KAFKA_BROKERS:114.55.58.251:9092} compression-type: zstd linger-ms: 20 @@ -568,9 +571,8 @@ lingniu: gb32960: enabled: false sink: - mq: + kafka: enabled: ${KAFKA_ENABLED:true} - type: kafka bootstrap-servers: ${KAFKA_BROKERS:114.55.58.251:9092} topics: realtime: ${KAFKA_TOPIC_GB32960_EVENT:vehicle.event.gb32960.v1} @@ -666,9 +668,8 @@ lingniu: gb32960: enabled: false sink: - mq: + kafka: enabled: ${KAFKA_ENABLED:true} - type: kafka bootstrap-servers: ${KAFKA_BROKERS:114.55.58.251:9092} topics: realtime: ${KAFKA_TOPIC_GB32960_EVENT:vehicle.event.gb32960.v1} @@ -754,7 +755,7 @@ import com.lingniu.ingest.eventfilestore.EventFileStore; import com.lingniu.ingest.protocol.gb32960.config.Gb32960AutoConfiguration; import com.lingniu.ingest.protocol.gb32960.inbound.Gb32960NettyServer; import com.lingniu.ingest.sink.archive.ArchiveStore; -import com.lingniu.ingest.sink.mq.KafkaEventSink; +import com.lingniu.ingest.sink.kafka.KafkaEventSink; import org.junit.jupiter.api.Test; import org.springframework.boot.autoconfigure.AutoConfigurations; import org.springframework.boot.test.context.runner.ApplicationContextRunner; @@ -770,7 +771,7 @@ class Gb32960IngestAppCompositionTest { .withPropertyValues( "lingniu.ingest.gb32960.enabled=true", "lingniu.ingest.gb32960.port=0", - "lingniu.ingest.sink.mq.enabled=false", + "lingniu.ingest.sink.kafka.enabled=false", "lingniu.ingest.sink.archive.enabled=false", "lingniu.ingest.event-file-store.enabled=false", "lingniu.ingest.event-history.enabled=false", @@ -835,7 +836,7 @@ class VehicleHistoryAppCompositionTest { "lingniu.ingest.event-history.enabled=true", "lingniu.ingest.vehicle-state.enabled=false", "lingniu.ingest.vehicle-stat.enabled=false", - "lingniu.ingest.sink.mq.enabled=false") + "lingniu.ingest.sink.kafka.enabled=false") .withConfiguration(AutoConfigurations.of( com.lingniu.ingest.sink.archive.config.SinkArchiveAutoConfiguration.class, com.lingniu.ingest.eventfilestore.config.EventFileStoreAutoConfiguration.class, @@ -896,7 +897,7 @@ class VehicleAnalyticsAppCompositionTest { "lingniu.ingest.vehicle-state.enabled=false", "lingniu.ingest.vehicle-stat.enabled=true", "lingniu.ingest.vehicle-stat.file-path=" + tempDir.resolve("vehicle-stat"), - "lingniu.ingest.sink.mq.enabled=false") + "lingniu.ingest.sink.kafka.enabled=false") .withConfiguration(AutoConfigurations.of( com.lingniu.ingest.vehiclestat.config.VehicleStatAutoConfiguration.class, com.lingniu.ingest.vehiclestate.config.VehicleStateAutoConfiguration.class)) @@ -932,7 +933,7 @@ 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/SinkMqProperties.java` +- 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` @@ -947,7 +948,7 @@ Expected test names: ```java @Test void rawArchiveRoutesToVersionedGb32960RawTopic() { - SinkMqProperties.Topics topics = new SinkMqProperties.Topics(); + KafkaSinkProperties.Topics topics = new KafkaSinkProperties.Topics(); topics.setRawArchive("vehicle.raw.gb32960.v1"); TopicRouter router = new TopicRouter(topics); @@ -970,7 +971,7 @@ Also add a normalized realtime event test: ```java @Test void realtimeRoutesToVersionedGb32960EventTopic() { - SinkMqProperties.Topics topics = new SinkMqProperties.Topics(); + KafkaSinkProperties.Topics topics = new KafkaSinkProperties.Topics(); topics.setRealtime("vehicle.event.gb32960.v1"); TopicRouter router = new TopicRouter(topics); @@ -1000,9 +1001,9 @@ mvn -pl :sink-mq -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. -- [ ] **Step 3: Update `SinkMqProperties.Topics` defaults** +- [ ] **Step 3: Update `KafkaSinkProperties.Topics` defaults** -In `SinkMqProperties.Topics`, change defaults: +In `KafkaSinkProperties.Topics`, change defaults: ```java private String realtime = "vehicle.event.gb32960.v1"; @@ -1040,7 +1041,7 @@ Expected: PASS. - [ ] **Step 6: Commit** ```bash -git add modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqProperties.java \ +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 \ @@ -1353,8 +1354,8 @@ void gb32960IngestDefaultsOnlyEnableProtocolAndKafkaProducer() { assertThat(props.getProperty("spring.application.name")).isEqualTo("gb32960-ingest-app"); assertThat(props.getProperty("lingniu.ingest.gb32960.enabled")).isEqualTo("true"); assertThat(props.getProperty("lingniu.ingest.gb32960.port")).isEqualTo("${GB32960_PORT:32960}"); - assertThat(props.getProperty("lingniu.ingest.sink.mq.enabled")).isEqualTo("${KAFKA_ENABLED:true}"); - assertThat(props.getProperty("lingniu.ingest.sink.mq.consumer.enabled")).isEqualTo("false"); + assertThat(props.getProperty("lingniu.ingest.sink.kafka.enabled")).isEqualTo("${KAFKA_ENABLED:true}"); + assertThat(props.getProperty("lingniu.ingest.sink.kafka.consumer.enabled")).isEqualTo("false"); assertThat(props.getProperty("lingniu.ingest.event-file-store.enabled")).isEqualTo("false"); assertThat(props.getProperty("lingniu.ingest.event-history.enabled")).isEqualTo("false"); assertThat(props.getProperty("lingniu.ingest.vehicle-state.enabled")).isEqualTo("false"); @@ -1378,7 +1379,7 @@ void historyDefaultsEnableStorageAndHistoryConsumerOnly() { assertThat(props.getProperty("lingniu.ingest.event-history.enabled")).isEqualTo("true"); assertThat(props.getProperty("lingniu.ingest.vehicle-state.enabled")).isEqualTo("false"); assertThat(props.getProperty("lingniu.ingest.vehicle-stat.enabled")).isEqualTo("false"); - assertThat(props.getProperty("lingniu.ingest.sink.mq.consumer.bindings.eventHistoryEnvelopeConsumerProcessor.group-id")) + assertThat(props.getProperty("lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryEnvelopeConsumerProcessor.group-id")) .isEqualTo("${KAFKA_GROUP_HISTORY:vehicle-history}"); } ``` @@ -1398,7 +1399,7 @@ void analyticsDefaultsEnableStatConsumerAndDisableProtocolAndHistoryStorage() { assertThat(props.getProperty("lingniu.ingest.event-file-store.enabled")).isEqualTo("false"); assertThat(props.getProperty("lingniu.ingest.event-history.enabled")).isEqualTo("false"); assertThat(props.getProperty("lingniu.ingest.vehicle-stat.enabled")).isEqualTo("${VEHICLE_STAT_ENABLED:true}"); - assertThat(props.getProperty("lingniu.ingest.sink.mq.consumer.bindings.vehicleStatEnvelopeConsumerProcessor.group-id")) + assertThat(props.getProperty("lingniu.ingest.sink.kafka.consumer.bindings.vehicleStatEnvelopeConsumerProcessor.group-id")) .isEqualTo("${KAFKA_GROUP_STAT:vehicle-stat}"); } ``` 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 a0fe0fe4..36e82acf 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 @@ -1269,12 +1269,12 @@ git commit -m "feat: add raw archive store contract" Create `VehicleEnvelopeProtoCompatibilityTest.java`: ```java -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; -import com.lingniu.ingest.sink.mq.proto.DecodedFactPayload; -import com.lingniu.ingest.sink.mq.proto.ParseStatusProto; -import com.lingniu.ingest.sink.mq.proto.RawFrameFactPayload; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.DecodedFactPayload; +import com.lingniu.ingest.sink.kafka.proto.ParseStatusProto; +import com.lingniu.ingest.sink.kafka.proto.RawFrameFactPayload; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import org.junit.jupiter.api.Test; import static org.assertj.core.api.Assertions.assertThat; 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 845bd6c4..d6d386fb 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 @@ -133,7 +133,7 @@ Use versioned topics: - `vehicle.event.gb32960.v1` - `vehicle.dlq.gb32960.v1` -The current topic names under `lingniu.ingest.sink.mq.topics` can remain temporarily for compatibility, but the split apps should converge on versioned topic names. +The current topic names under `lingniu.ingest.sink.kafka.topics` can remain temporarily for compatibility, but the split apps should converge on versioned topic names. ### Keys diff --git a/docs/target-architecture.md b/docs/target-architecture.md index 758a6416..7148e748 100644 --- a/docs/target-architecture.md +++ b/docs/target-architecture.md @@ -71,7 +71,7 @@ repositories, and statistic calculators. ### Event Contract -`sink-mq` owns the Kafka protobuf envelope and consumer plumbing. +The Kafka sink owns the protobuf envelope and consumer plumbing. Consumers should use `EnvelopeConsumerProcessor` so invalid protobuf, skipped envelopes, and downstream failures are converted to structured results and DLQ records instead of blocking a partition. diff --git a/modules/apps/gb32960-ingest-app/src/main/resources/application.yml b/modules/apps/gb32960-ingest-app/src/main/resources/application.yml index ba28800a..49397d62 100644 --- a/modules/apps/gb32960-ingest-app/src/main/resources/application.yml +++ b/modules/apps/gb32960-ingest-app/src/main/resources/application.yml @@ -99,9 +99,8 @@ lingniu: table-name: ${VEHICLE_IDENTITY_MYSQL_TABLE:vehicle_identity_bindings} initialize-schema: ${VEHICLE_IDENTITY_MYSQL_INITIALIZE_SCHEMA:true} sink: - mq: + kafka: enabled: ${KAFKA_ENABLED:true} - type: kafka bootstrap-servers: ${KAFKA_BROKERS:114.55.58.251:9092} compression-type: zstd linger-ms: 20 diff --git a/modules/apps/gb32960-ingest-app/src/test/java/com/lingniu/ingest/gb32960app/Gb32960IngestAppCompositionTest.java b/modules/apps/gb32960-ingest-app/src/test/java/com/lingniu/ingest/gb32960app/Gb32960IngestAppCompositionTest.java index b47ca5e4..4a3f10cd 100644 --- a/modules/apps/gb32960-ingest-app/src/test/java/com/lingniu/ingest/gb32960app/Gb32960IngestAppCompositionTest.java +++ b/modules/apps/gb32960-ingest-app/src/test/java/com/lingniu/ingest/gb32960app/Gb32960IngestAppCompositionTest.java @@ -11,9 +11,9 @@ import com.lingniu.ingest.session.config.SessionCoreAutoConfiguration; import com.lingniu.ingest.sink.archive.ArchiveStore; import com.lingniu.ingest.sink.archive.RawArchiveEventSink; import com.lingniu.ingest.sink.archive.config.SinkArchiveAutoConfiguration; -import com.lingniu.ingest.sink.mq.KafkaEventSink; -import com.lingniu.ingest.sink.mq.KafkaEnvelopeDeadLetterSink; -import com.lingniu.ingest.sink.mq.SinkMqAutoConfiguration; +import com.lingniu.ingest.sink.kafka.KafkaEventSink; +import com.lingniu.ingest.sink.kafka.KafkaEnvelopeDeadLetterSink; +import com.lingniu.ingest.sink.kafka.KafkaSinkAutoConfiguration; import org.apache.kafka.clients.producer.KafkaProducer; import org.junit.jupiter.api.Test; import org.springframework.boot.autoconfigure.AutoConfigurations; @@ -31,7 +31,7 @@ class Gb32960IngestAppCompositionTest { IngestCoreAutoConfiguration.class, SessionCoreAutoConfiguration.class, VehicleIdentityAutoConfiguration.class, - SinkMqAutoConfiguration.class, + KafkaSinkAutoConfiguration.class, SinkArchiveAutoConfiguration.class, Gb32960AutoConfiguration.class)) .withAllowBeanDefinitionOverriding(true) @@ -43,10 +43,9 @@ class Gb32960IngestAppCompositionTest { "lingniu.ingest.gb32960.port=0", "lingniu.ingest.identity.store=mysql", "lingniu.ingest.identity.mysql.initialize-schema=false", - "lingniu.ingest.sink.mq.enabled=true", - "lingniu.ingest.sink.mq.type=kafka", - "lingniu.ingest.sink.mq.bootstrap-servers=localhost:9092", - "lingniu.ingest.sink.mq.consumer.enabled=false", + "lingniu.ingest.sink.kafka.enabled=true", + "lingniu.ingest.sink.kafka.bootstrap-servers=localhost:9092", + "lingniu.ingest.sink.kafka.consumer.enabled=false", "lingniu.ingest.sink.archive.enabled=true", "lingniu.ingest.sink.archive.path=target/test-archive-gb32960", "lingniu.ingest.event-history.enabled=false"); diff --git a/modules/apps/gb32960-ingest-app/src/test/java/com/lingniu/ingest/gb32960app/Gb32960IngestAppDefaultsTest.java b/modules/apps/gb32960-ingest-app/src/test/java/com/lingniu/ingest/gb32960app/Gb32960IngestAppDefaultsTest.java index 8c4908b8..1f9904df 100644 --- a/modules/apps/gb32960-ingest-app/src/test/java/com/lingniu/ingest/gb32960app/Gb32960IngestAppDefaultsTest.java +++ b/modules/apps/gb32960-ingest-app/src/test/java/com/lingniu/ingest/gb32960app/Gb32960IngestAppDefaultsTest.java @@ -15,6 +15,8 @@ class Gb32960IngestAppDefaultsTest { @Test void applicationDefaultsKeepGb32960IngestAsTcpKafkaProducerOnly() throws IOException { Properties properties = applicationProperties(); + String legacyKafkaPrefix = "lingniu.ingest.sink." + "mq."; + String kafkaTypeProperty = "lingniu.ingest.sink.kafka." + "type"; assertThat(properties) .containsEntry("spring.application.name", "gb32960-ingest-app") @@ -40,8 +42,8 @@ class Gb32960IngestAppDefaultsTest { "Hyundai") .containsEntry("lingniu.ingest.gb32960.vendor-extensions[0].match.platform-accounts[1]", "YueJin") - .containsEntry("lingniu.ingest.sink.mq.enabled", "${KAFKA_ENABLED:true}") - .containsEntry("lingniu.ingest.sink.mq.consumer.enabled", false) + .containsEntry("lingniu.ingest.sink.kafka.enabled", "${KAFKA_ENABLED:true}") + .containsEntry("lingniu.ingest.sink.kafka.consumer.enabled", false) .containsEntry("lingniu.ingest.sink.archive.enabled", "${SINK_ARCHIVE_ENABLED:true}") .containsEntry("lingniu.ingest.sink.archive.path", "${SINK_ARCHIVE_PATH:./archive/}") .containsEntry("lingniu.ingest.session.store", "${SESSION_STORE:redis}") @@ -55,7 +57,9 @@ class Gb32960IngestAppDefaultsTest { assertThat(properties.stringPropertyNames()) .noneMatch(name -> name.startsWith("lingniu.ingest.event-file-store.")) .noneMatch(name -> name.startsWith("lingniu.ingest.vehicle-state.")) - .noneMatch(name -> name.startsWith("lingniu.ingest.vehicle-stat.")); + .noneMatch(name -> name.startsWith("lingniu.ingest.vehicle-stat.")) + .noneMatch(name -> name.startsWith(legacyKafkaPrefix)) + .noneMatch(name -> name.equals(kafkaTypeProperty)); assertThat(applicationYaml()) .doesNotContain("event-file-store:") .doesNotContain("vehicle-state:") diff --git a/modules/apps/jt808-ingest-app/src/main/resources/application.yml b/modules/apps/jt808-ingest-app/src/main/resources/application.yml index bf4a2719..9622a619 100644 --- a/modules/apps/jt808-ingest-app/src/main/resources/application.yml +++ b/modules/apps/jt808-ingest-app/src/main/resources/application.yml @@ -63,9 +63,8 @@ lingniu: table-name: ${VEHICLE_IDENTITY_MYSQL_TABLE:vehicle_identity_bindings} initialize-schema: ${VEHICLE_IDENTITY_MYSQL_INITIALIZE_SCHEMA:true} sink: - mq: + kafka: enabled: ${KAFKA_ENABLED:true} - type: kafka bootstrap-servers: ${KAFKA_BROKERS:114.55.58.251:9092} compression-type: zstd linger-ms: 20 diff --git a/modules/apps/jt808-ingest-app/src/test/java/com/lingniu/ingest/jt808app/Jt808IngestAppCompositionTest.java b/modules/apps/jt808-ingest-app/src/test/java/com/lingniu/ingest/jt808app/Jt808IngestAppCompositionTest.java index e00ee080..e47438b8 100644 --- a/modules/apps/jt808-ingest-app/src/test/java/com/lingniu/ingest/jt808app/Jt808IngestAppCompositionTest.java +++ b/modules/apps/jt808-ingest-app/src/test/java/com/lingniu/ingest/jt808app/Jt808IngestAppCompositionTest.java @@ -11,9 +11,9 @@ import com.lingniu.ingest.session.config.SessionCoreAutoConfiguration; import com.lingniu.ingest.sink.archive.ArchiveStore; import com.lingniu.ingest.sink.archive.RawArchiveEventSink; import com.lingniu.ingest.sink.archive.config.SinkArchiveAutoConfiguration; -import com.lingniu.ingest.sink.mq.KafkaEnvelopeDeadLetterSink; -import com.lingniu.ingest.sink.mq.KafkaEventSink; -import com.lingniu.ingest.sink.mq.SinkMqAutoConfiguration; +import com.lingniu.ingest.sink.kafka.KafkaEnvelopeDeadLetterSink; +import com.lingniu.ingest.sink.kafka.KafkaEventSink; +import com.lingniu.ingest.sink.kafka.KafkaSinkAutoConfiguration; import org.apache.kafka.clients.producer.KafkaProducer; import org.junit.jupiter.api.Test; import org.springframework.boot.autoconfigure.AutoConfigurations; @@ -31,7 +31,7 @@ class Jt808IngestAppCompositionTest { IngestCoreAutoConfiguration.class, SessionCoreAutoConfiguration.class, VehicleIdentityAutoConfiguration.class, - SinkMqAutoConfiguration.class, + KafkaSinkAutoConfiguration.class, SinkArchiveAutoConfiguration.class, Jt808AutoConfiguration.class)) .withAllowBeanDefinitionOverriding(true) @@ -42,10 +42,9 @@ class Jt808IngestAppCompositionTest { "lingniu.ingest.jt808.port=0", "lingniu.ingest.identity.store=mysql", "lingniu.ingest.identity.mysql.initialize-schema=false", - "lingniu.ingest.sink.mq.enabled=true", - "lingniu.ingest.sink.mq.type=kafka", - "lingniu.ingest.sink.mq.bootstrap-servers=localhost:9092", - "lingniu.ingest.sink.mq.consumer.enabled=false", + "lingniu.ingest.sink.kafka.enabled=true", + "lingniu.ingest.sink.kafka.bootstrap-servers=localhost:9092", + "lingniu.ingest.sink.kafka.consumer.enabled=false", "lingniu.ingest.sink.archive.enabled=true", "lingniu.ingest.sink.archive.path=target/test-archive-jt808", "lingniu.ingest.event-history.enabled=false"); diff --git a/modules/apps/jt808-ingest-app/src/test/java/com/lingniu/ingest/jt808app/Jt808IngestAppDefaultsTest.java b/modules/apps/jt808-ingest-app/src/test/java/com/lingniu/ingest/jt808app/Jt808IngestAppDefaultsTest.java index c651403c..c1e20c39 100644 --- a/modules/apps/jt808-ingest-app/src/test/java/com/lingniu/ingest/jt808app/Jt808IngestAppDefaultsTest.java +++ b/modules/apps/jt808-ingest-app/src/test/java/com/lingniu/ingest/jt808app/Jt808IngestAppDefaultsTest.java @@ -28,11 +28,11 @@ class Jt808IngestAppDefaultsTest { .containsEntry("server.port", "${HTTP_PORT:20400}") .containsEntry("lingniu.ingest.jt808.enabled", true) .containsEntry("lingniu.ingest.jt808.port", "${JT808_PORT:808}") - .containsEntry("lingniu.ingest.sink.mq.enabled", "${KAFKA_ENABLED:true}") - .containsEntry("lingniu.ingest.sink.mq.consumer.enabled", false) - .containsEntry("lingniu.ingest.sink.mq.topics.realtime", "${KAFKA_TOPIC_JT808_EVENT:vehicle.event.jt808.v1}") - .containsEntry("lingniu.ingest.sink.mq.topics.raw-archive", "${KAFKA_TOPIC_JT808_RAW:vehicle.raw.jt808.v1}") - .containsEntry("lingniu.ingest.sink.mq.topics.dlq", "${KAFKA_TOPIC_JT808_DLQ:vehicle.dlq.jt808.v1}") + .containsEntry("lingniu.ingest.sink.kafka.enabled", "${KAFKA_ENABLED:true}") + .containsEntry("lingniu.ingest.sink.kafka.consumer.enabled", false) + .containsEntry("lingniu.ingest.sink.kafka.topics.realtime", "${KAFKA_TOPIC_JT808_EVENT:vehicle.event.jt808.v1}") + .containsEntry("lingniu.ingest.sink.kafka.topics.raw-archive", "${KAFKA_TOPIC_JT808_RAW:vehicle.raw.jt808.v1}") + .containsEntry("lingniu.ingest.sink.kafka.topics.dlq", "${KAFKA_TOPIC_JT808_DLQ:vehicle.dlq.jt808.v1}") .containsEntry("lingniu.ingest.session.store", "${SESSION_STORE:redis}") .containsEntry("lingniu.ingest.identity.store", "${VEHICLE_IDENTITY_STORE:mysql}") .containsEntry("lingniu.ingest.identity.mysql.table-name", diff --git a/modules/apps/vehicle-analytics-app/src/main/resources/application.yml b/modules/apps/vehicle-analytics-app/src/main/resources/application.yml index 94ecf28e..482bb1fc 100644 --- a/modules/apps/vehicle-analytics-app/src/main/resources/application.yml +++ b/modules/apps/vehicle-analytics-app/src/main/resources/application.yml @@ -37,9 +37,8 @@ lingniu: gb32960: enabled: false sink: - mq: + kafka: enabled: ${KAFKA_ENABLED:true} - type: kafka bootstrap-servers: ${KAFKA_BROKERS:114.55.58.251:9092} topics: realtime: ${KAFKA_TOPIC_JT808_EVENT:vehicle.event.jt808.v1} diff --git a/modules/apps/vehicle-analytics-app/src/test/java/com/lingniu/ingest/analyticsapp/VehicleAnalyticsAppCompositionTest.java b/modules/apps/vehicle-analytics-app/src/test/java/com/lingniu/ingest/analyticsapp/VehicleAnalyticsAppCompositionTest.java index a31b5c42..133d7999 100644 --- a/modules/apps/vehicle-analytics-app/src/test/java/com/lingniu/ingest/analyticsapp/VehicleAnalyticsAppCompositionTest.java +++ b/modules/apps/vehicle-analytics-app/src/test/java/com/lingniu/ingest/analyticsapp/VehicleAnalyticsAppCompositionTest.java @@ -2,9 +2,9 @@ package com.lingniu.ingest.analyticsapp; import com.fasterxml.jackson.databind.ObjectMapper; import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor; -import com.lingniu.ingest.sink.mq.KafkaEnvelopeDeadLetterSink; -import com.lingniu.ingest.sink.mq.KafkaEventSink; -import com.lingniu.ingest.sink.mq.SinkMqAutoConfiguration; +import com.lingniu.ingest.sink.kafka.KafkaEnvelopeDeadLetterSink; +import com.lingniu.ingest.sink.kafka.KafkaEventSink; +import com.lingniu.ingest.sink.kafka.KafkaSinkAutoConfiguration; import com.lingniu.ingest.vehiclestat.JdbcVehicleStatMetricRepository; import com.lingniu.ingest.vehiclestat.VehicleStatController; import com.lingniu.ingest.vehiclestat.VehicleStatEnvelopeIngestor; @@ -28,17 +28,16 @@ class VehicleAnalyticsAppCompositionTest { void createsStatsBeansWithoutProtocolListenerOrEventFileStore() { new ApplicationContextRunner() .withConfiguration(AutoConfigurations.of( - SinkMqAutoConfiguration.class, + KafkaSinkAutoConfiguration.class, VehicleStatAutoConfiguration.class)) .withAllowBeanDefinitionOverriding(true) .withBean("kafkaProducer", KafkaProducer.class, VehicleAnalyticsAppCompositionTest::kafkaProducer) .withBean(JdbcTemplate.class, () -> mock(JdbcTemplate.class)) .withBean(ObjectMapper.class, ObjectMapper::new) .withPropertyValues( - "lingniu.ingest.sink.mq.enabled=true", - "lingniu.ingest.sink.mq.type=kafka", - "lingniu.ingest.sink.mq.bootstrap-servers=localhost:9092", - "lingniu.ingest.sink.mq.consumer.enabled=false", + "lingniu.ingest.sink.kafka.enabled=true", + "lingniu.ingest.sink.kafka.bootstrap-servers=localhost:9092", + "lingniu.ingest.sink.kafka.consumer.enabled=false", "lingniu.ingest.vehicle-stat.enabled=true", "lingniu.ingest.vehicle-stat.jt808.enabled=true", "lingniu.ingest.event-history.enabled=false", diff --git a/modules/apps/vehicle-analytics-app/src/test/java/com/lingniu/ingest/analyticsapp/VehicleAnalyticsAppDefaultsTest.java b/modules/apps/vehicle-analytics-app/src/test/java/com/lingniu/ingest/analyticsapp/VehicleAnalyticsAppDefaultsTest.java index 0ba809bb..d60bb4b1 100644 --- a/modules/apps/vehicle-analytics-app/src/test/java/com/lingniu/ingest/analyticsapp/VehicleAnalyticsAppDefaultsTest.java +++ b/modules/apps/vehicle-analytics-app/src/test/java/com/lingniu/ingest/analyticsapp/VehicleAnalyticsAppDefaultsTest.java @@ -29,12 +29,12 @@ class VehicleAnalyticsAppDefaultsTest { .containsEntry( "lingniu.ingest.vehicle-stat.jt808.enabled", "${VEHICLE_STAT_JT808_MILEAGE_ENABLED:true}") - .containsEntry("lingniu.ingest.sink.mq.consumer.enabled", "${KAFKA_CONSUMER_ENABLED:true}") + .containsEntry("lingniu.ingest.sink.kafka.consumer.enabled", "${KAFKA_CONSUMER_ENABLED:true}") .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.vehicleStatEnvelopeConsumerProcessor.enabled", + "lingniu.ingest.sink.kafka.consumer.bindings.vehicleStatEnvelopeConsumerProcessor.enabled", "${VEHICLE_STAT_ENABLED:true}") .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.vehicleStatEnvelopeConsumerProcessor.group-id", + "lingniu.ingest.sink.kafka.consumer.bindings.vehicleStatEnvelopeConsumerProcessor.group-id", "${KAFKA_GROUP_STAT:vehicle-stat}") .containsEntry("management.endpoints.web.exposure.include", "health,info,metrics,prometheus"); diff --git a/modules/apps/vehicle-history-app/src/main/java/com/lingniu/ingest/historyapp/VehicleHistoryKafkaConsumerConfiguration.java b/modules/apps/vehicle-history-app/src/main/java/com/lingniu/ingest/historyapp/VehicleHistoryKafkaConsumerConfiguration.java index a46fa39f..5884c23f 100644 --- a/modules/apps/vehicle-history-app/src/main/java/com/lingniu/ingest/historyapp/VehicleHistoryKafkaConsumerConfiguration.java +++ b/modules/apps/vehicle-history-app/src/main/java/com/lingniu/ingest/historyapp/VehicleHistoryKafkaConsumerConfiguration.java @@ -1,10 +1,10 @@ package com.lingniu.ingest.historyapp; import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor; -import com.lingniu.ingest.sink.mq.KafkaEnvelopeConsumerFactory; -import com.lingniu.ingest.sink.mq.KafkaEnvelopeConsumerRunner; -import com.lingniu.ingest.sink.mq.KafkaEnvelopeConsumerWorker; -import com.lingniu.ingest.sink.mq.SinkMqProperties; +import com.lingniu.ingest.sink.kafka.KafkaEnvelopeConsumerFactory; +import com.lingniu.ingest.sink.kafka.KafkaEnvelopeConsumerRunner; +import com.lingniu.ingest.sink.kafka.KafkaEnvelopeConsumerWorker; +import com.lingniu.ingest.sink.kafka.KafkaSinkProperties; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; @@ -20,11 +20,11 @@ public class VehicleHistoryKafkaConsumerConfiguration { @Bean @ConditionalOnMissingBean - @ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq.consumer", name = "enabled", havingValue = "true") + @ConditionalOnProperty(prefix = "lingniu.ingest.sink.kafka.consumer", name = "enabled", havingValue = "true") public KafkaEnvelopeConsumerRunner vehicleHistoryKafkaEnvelopeConsumerRunner( @Qualifier("eventHistoryEnvelopeConsumerProcessor") EnvelopeConsumerProcessor processor, @Qualifier("eventHistoryRawEnvelopeConsumerProcessor") EnvelopeConsumerProcessor rawProcessor, - SinkMqProperties props) { + KafkaSinkProperties props) { List workers = new KafkaEnvelopeConsumerFactory().createWorkers( Map.of( "eventHistoryGb32960EnvelopeConsumerProcessor", processor, diff --git a/modules/apps/vehicle-history-app/src/main/resources/application.yml b/modules/apps/vehicle-history-app/src/main/resources/application.yml index c2dea2f1..859bf837 100644 --- a/modules/apps/vehicle-history-app/src/main/resources/application.yml +++ b/modules/apps/vehicle-history-app/src/main/resources/application.yml @@ -34,9 +34,8 @@ lingniu: server: enabled: false sink: - mq: + kafka: enabled: ${KAFKA_ENABLED:true} - type: kafka bootstrap-servers: ${KAFKA_BROKERS:114.55.58.251:9092} topics: realtime: ${KAFKA_TOPIC_GB32960_EVENT:vehicle.event.gb32960.v1} diff --git a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppCompositionTest.java b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppCompositionTest.java index 6fed29bc..fa4ab4d7 100644 --- a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppCompositionTest.java +++ b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppCompositionTest.java @@ -15,15 +15,15 @@ import com.lingniu.ingest.eventhistory.TelemetryFieldHistoryController; import com.lingniu.ingest.protocol.gb32960.codec.Gb32960MessageDecoder; import com.lingniu.ingest.protocol.gb32960.config.Gb32960AutoConfiguration; import com.lingniu.ingest.protocol.gb32960.inbound.Gb32960NettyServer; -import com.lingniu.ingest.sink.mq.KafkaEnvelopeDeadLetterSink; -import com.lingniu.ingest.sink.mq.KafkaEventSink; -import com.lingniu.ingest.sink.mq.KafkaEnvelopeConsumerRunner; -import com.lingniu.ingest.sink.mq.SinkMqProperties; -import com.lingniu.ingest.sink.mq.proto.ParseStatusProto; -import com.lingniu.ingest.sink.mq.proto.RawArchiveRef; -import com.lingniu.ingest.sink.mq.proto.RawFrameFactPayload; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; -import com.lingniu.ingest.sink.mq.SinkMqAutoConfiguration; +import com.lingniu.ingest.sink.kafka.KafkaEnvelopeDeadLetterSink; +import com.lingniu.ingest.sink.kafka.KafkaEventSink; +import com.lingniu.ingest.sink.kafka.KafkaEnvelopeConsumerRunner; +import com.lingniu.ingest.sink.kafka.KafkaSinkProperties; +import com.lingniu.ingest.sink.kafka.proto.ParseStatusProto; +import com.lingniu.ingest.sink.kafka.proto.RawArchiveRef; +import com.lingniu.ingest.sink.kafka.proto.RawFrameFactPayload; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.KafkaSinkAutoConfiguration; import com.lingniu.ingest.tdenginehistory.TdengineHistorySchema; import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter; import com.lingniu.ingest.tdenginehistory.TdengineRawFrameRow; @@ -75,8 +75,8 @@ class VehicleHistoryAppCompositionTest { () -> processor("event-history")) .withBean("eventHistoryRawEnvelopeConsumerProcessor", EnvelopeConsumerProcessor.class, () -> processor("event-history-raw")) - .withBean(SinkMqProperties.class, VehicleHistoryAppCompositionTest::genericOnlyHistoryConsumerProps) - .withPropertyValues("lingniu.ingest.sink.mq.consumer.enabled=true") + .withBean(KafkaSinkProperties.class, VehicleHistoryAppCompositionTest::genericOnlyHistoryConsumerProps) + .withPropertyValues("lingniu.ingest.sink.kafka.consumer.enabled=true") .run(context -> { assertThat(context).hasFailed(); assertThat(context.getStartupFailure()) @@ -89,7 +89,7 @@ class VehicleHistoryAppCompositionTest { new ApplicationContextRunner() .withConfiguration(AutoConfigurations.of( TdengineHistoryAutoConfiguration.class, - SinkMqAutoConfiguration.class, + KafkaSinkAutoConfiguration.class, Gb32960AutoConfiguration.class, EventHistoryAutoConfiguration.class)) .withUserConfiguration(VehicleHistoryKafkaConsumerConfiguration.class) @@ -101,10 +101,9 @@ class VehicleHistoryAppCompositionTest { "lingniu.ingest.event-history.enabled=true", "lingniu.ingest.gb32960.enabled=true", "lingniu.ingest.gb32960.server.enabled=false", - "lingniu.ingest.sink.mq.enabled=true", - "lingniu.ingest.sink.mq.type=kafka", - "lingniu.ingest.sink.mq.bootstrap-servers=localhost:9092", - "lingniu.ingest.sink.mq.consumer.enabled=false", + "lingniu.ingest.sink.kafka.enabled=true", + "lingniu.ingest.sink.kafka.bootstrap-servers=localhost:9092", + "lingniu.ingest.sink.kafka.consumer.enabled=false", "lingniu.ingest.vehicle-state.enabled=false", "lingniu.ingest.vehicle-stat.enabled=false") .run(context -> { @@ -137,7 +136,7 @@ class VehicleHistoryAppCompositionTest { .withBean(EnvelopeDeadLetterSink.class, () -> mock(EnvelopeDeadLetterSink.class)) .withPropertyValues( "lingniu.ingest.event-history.enabled=true", - "lingniu.ingest.sink.mq.consumer.enabled=false") + "lingniu.ingest.sink.kafka.consumer.enabled=false") .run(context -> { assertThat(context).hasSingleBean(EventHistoryEnvelopeIngestor.class); assertThat(context).doesNotHaveBean("eventFileStore"); @@ -149,7 +148,7 @@ class VehicleHistoryAppCompositionTest { new ApplicationContextRunner() .withConfiguration(AutoConfigurations.of( TdengineHistoryAutoConfiguration.class, - SinkMqAutoConfiguration.class, + KafkaSinkAutoConfiguration.class, EventHistoryAutoConfiguration.class)) .withUserConfiguration(VehicleHistoryKafkaConsumerConfiguration.class) .withAllowBeanDefinitionOverriding(true) @@ -158,13 +157,12 @@ class VehicleHistoryAppCompositionTest { "lingniu.ingest.tdengine-history.enabled=true", "lingniu.ingest.tdengine-history.database=vehicle_history_test", "lingniu.ingest.event-history.enabled=true", - "lingniu.ingest.sink.mq.enabled=true", - "lingniu.ingest.sink.mq.type=kafka", - "lingniu.ingest.sink.mq.bootstrap-servers=localhost:9092", - "lingniu.ingest.sink.mq.consumer.enabled=true", - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryGb32960EnvelopeConsumerProcessor.enabled=true", - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryGb32960EnvelopeConsumerProcessor.group-id=history-gb32960", - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryGb32960EnvelopeConsumerProcessor.topics[0]=vehicle.event.gb32960.v1") + "lingniu.ingest.sink.kafka.enabled=true", + "lingniu.ingest.sink.kafka.bootstrap-servers=localhost:9092", + "lingniu.ingest.sink.kafka.consumer.enabled=true", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryGb32960EnvelopeConsumerProcessor.enabled=true", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryGb32960EnvelopeConsumerProcessor.group-id=history-gb32960", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryGb32960EnvelopeConsumerProcessor.topics[0]=vehicle.event.gb32960.v1") .run(context -> { assertThat(context).hasSingleBean(EventHistoryEnvelopeIngestor.class); assertThat(context.getBeanNamesForType(EnvelopeConsumerProcessor.class)) @@ -186,7 +184,7 @@ class VehicleHistoryAppCompositionTest { .withBean(EnvelopeDeadLetterSink.class, () -> mock(EnvelopeDeadLetterSink.class)) .withPropertyValues( "lingniu.ingest.event-history.enabled=true", - "lingniu.ingest.sink.mq.consumer.enabled=false") + "lingniu.ingest.sink.kafka.consumer.enabled=false") .run(context -> { EventHistoryEnvelopeIngestor ingestor = context.getBean(EventHistoryEnvelopeIngestor.class); @@ -211,8 +209,8 @@ class VehicleHistoryAppCompositionTest { ignored -> {}); } - private static SinkMqProperties genericOnlyHistoryConsumerProps() { - SinkMqProperties props = new SinkMqProperties(); + private static KafkaSinkProperties genericOnlyHistoryConsumerProps() { + KafkaSinkProperties props = new KafkaSinkProperties(); props.setBootstrapServers("localhost:9092"); props.getConsumer().setEnabled(true); props.getConsumer().setBindings(Map.of( @@ -221,8 +219,8 @@ class VehicleHistoryAppCompositionTest { return props; } - private static SinkMqProperties.Binding binding(String groupId, String topic) { - SinkMqProperties.Binding binding = new SinkMqProperties.Binding(); + private static KafkaSinkProperties.Binding binding(String groupId, String topic) { + KafkaSinkProperties.Binding binding = new KafkaSinkProperties.Binding(); binding.setGroupId(groupId); binding.setTopics(List.of(topic)); return binding; diff --git a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppDefaultsTest.java b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppDefaultsTest.java index 55d36dfa..1bb76302 100644 --- a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppDefaultsTest.java +++ b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppDefaultsTest.java @@ -44,65 +44,65 @@ class VehicleHistoryAppDefaultsTest { .containsEntry("lingniu.ingest.event-history.enabled", true) .containsEntry("lingniu.ingest.vehicle-state.enabled", false) .containsEntry("lingniu.ingest.vehicle-stat.enabled", false) - .containsEntry("lingniu.ingest.sink.mq.consumer.enabled", "${KAFKA_CONSUMER_ENABLED:true}") - .containsEntry("lingniu.ingest.sink.mq.consumer.max-poll-records", "${KAFKA_CONSUMER_MAX_POLL_RECORDS:2000}") - .containsEntry("lingniu.ingest.sink.mq.consumer.concurrency", "${KAFKA_CONSUMER_CONCURRENCY:3}") + .containsEntry("lingniu.ingest.sink.kafka.consumer.enabled", "${KAFKA_CONSUMER_ENABLED:true}") + .containsEntry("lingniu.ingest.sink.kafka.consumer.max-poll-records", "${KAFKA_CONSUMER_MAX_POLL_RECORDS:2000}") + .containsEntry("lingniu.ingest.sink.kafka.consumer.concurrency", "${KAFKA_CONSUMER_CONCURRENCY:3}") .containsEntry( - "lingniu.ingest.sink.mq.consumer.max-poll-interval-millis", + "lingniu.ingest.sink.kafka.consumer.max-poll-interval-millis", "${KAFKA_CONSUMER_MAX_POLL_INTERVAL_MS:1800000}") .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryGb32960EnvelopeConsumerProcessor.enabled", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryGb32960EnvelopeConsumerProcessor.enabled", true) .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryGb32960EnvelopeConsumerProcessor.group-id", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryGb32960EnvelopeConsumerProcessor.group-id", "${KAFKA_GROUP_HISTORY_GB32960_EVENT:${KAFKA_GROUP_HISTORY:vehicle-history}-gb32960-event}") .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryGb32960EnvelopeConsumerProcessor.topics[0]", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryGb32960EnvelopeConsumerProcessor.topics[0]", "${KAFKA_TOPIC_GB32960_EVENT:vehicle.event.gb32960.v1}") .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryJt808EnvelopeConsumerProcessor.enabled", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryJt808EnvelopeConsumerProcessor.enabled", true) .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryJt808EnvelopeConsumerProcessor.group-id", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryJt808EnvelopeConsumerProcessor.group-id", "${KAFKA_GROUP_HISTORY_JT808_EVENT:${KAFKA_GROUP_HISTORY:vehicle-history}-jt808-event}") .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryJt808EnvelopeConsumerProcessor.topics[0]", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryJt808EnvelopeConsumerProcessor.topics[0]", "${KAFKA_TOPIC_JT808_EVENT:vehicle.event.jt808.v1}") .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryYutongMqttEnvelopeConsumerProcessor.enabled", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryYutongMqttEnvelopeConsumerProcessor.enabled", true) .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryYutongMqttEnvelopeConsumerProcessor.group-id", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryYutongMqttEnvelopeConsumerProcessor.group-id", "${KAFKA_GROUP_HISTORY_YUTONG_MQTT_EVENT:${KAFKA_GROUP_HISTORY:vehicle-history}-yutong-mqtt-event}") .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryYutongMqttEnvelopeConsumerProcessor.topics[0]", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryYutongMqttEnvelopeConsumerProcessor.topics[0]", "${KAFKA_TOPIC_YUTONG_MQTT_EVENT:vehicle.event.mqtt-yutong.v1}") .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryGb32960RawEnvelopeConsumerProcessor.enabled", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryGb32960RawEnvelopeConsumerProcessor.enabled", true) .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryGb32960RawEnvelopeConsumerProcessor.group-id", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryGb32960RawEnvelopeConsumerProcessor.group-id", "${KAFKA_GROUP_HISTORY_GB32960_RAW:${KAFKA_GROUP_HISTORY:vehicle-history}-gb32960-raw}") .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryGb32960RawEnvelopeConsumerProcessor.topics[0]", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryGb32960RawEnvelopeConsumerProcessor.topics[0]", "${KAFKA_TOPIC_GB32960_RAW:vehicle.raw.gb32960.v1}") .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryJt808RawEnvelopeConsumerProcessor.enabled", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryJt808RawEnvelopeConsumerProcessor.enabled", true) .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryJt808RawEnvelopeConsumerProcessor.group-id", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryJt808RawEnvelopeConsumerProcessor.group-id", "${KAFKA_GROUP_HISTORY_JT808_RAW:${KAFKA_GROUP_HISTORY:vehicle-history}-jt808-raw}") .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryJt808RawEnvelopeConsumerProcessor.topics[0]", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryJt808RawEnvelopeConsumerProcessor.topics[0]", "${KAFKA_TOPIC_JT808_RAW:vehicle.raw.jt808.v1}") .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryYutongMqttRawEnvelopeConsumerProcessor.enabled", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryYutongMqttRawEnvelopeConsumerProcessor.enabled", true) .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryYutongMqttRawEnvelopeConsumerProcessor.group-id", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryYutongMqttRawEnvelopeConsumerProcessor.group-id", "${KAFKA_GROUP_HISTORY_YUTONG_MQTT_RAW:${KAFKA_GROUP_HISTORY:vehicle-history}-yutong-mqtt-raw}") .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryYutongMqttRawEnvelopeConsumerProcessor.topics[0]", + "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryYutongMqttRawEnvelopeConsumerProcessor.topics[0]", "${KAFKA_TOPIC_YUTONG_MQTT_RAW:vehicle.raw.mqtt-yutong.v1}") .containsEntry("management.endpoints.web.exposure.include", "health,info,metrics,prometheus"); diff --git a/modules/apps/xinda-push-app/src/main/resources/application.yml b/modules/apps/xinda-push-app/src/main/resources/application.yml index ee566266..775c4d84 100644 --- a/modules/apps/xinda-push-app/src/main/resources/application.yml +++ b/modules/apps/xinda-push-app/src/main/resources/application.yml @@ -53,20 +53,16 @@ lingniu: rate-limit: per-vin-qps: 50 identity: - store: ${VEHICLE_IDENTITY_STORE:file} - file: - path: ${VEHICLE_IDENTITY_FILE:./data/vehicle-identity.jsonl} + store: ${VEHICLE_IDENTITY_STORE:mysql} mysql: - table: ${VEHICLE_IDENTITY_MYSQL_TABLE:vehicle_identity_binding} - jdbc-url: ${VEHICLE_IDENTITY_MYSQL_JDBC_URL:} - username: ${VEHICLE_IDENTITY_MYSQL_USERNAME:} + jdbc-url: ${VEHICLE_IDENTITY_MYSQL_JDBC_URL:jdbc:mysql://127.0.0.1:3306/lingniu_vehicle?useUnicode=true&characterEncoding=utf8&useSSL=false&serverTimezone=Asia/Shanghai} + username: ${VEHICLE_IDENTITY_MYSQL_USERNAME:root} password: ${VEHICLE_IDENTITY_MYSQL_PASSWORD:} - driver-class-name: ${VEHICLE_IDENTITY_MYSQL_DRIVER_CLASS_NAME:com.mysql.cj.jdbc.Driver} - refresh-interval: ${VEHICLE_IDENTITY_MYSQL_REFRESH_INTERVAL:60s} + table-name: ${VEHICLE_IDENTITY_MYSQL_TABLE:vehicle_identity_bindings} + initialize-schema: ${VEHICLE_IDENTITY_MYSQL_INITIALIZE_SCHEMA:true} sink: - mq: + kafka: enabled: ${KAFKA_ENABLED:true} - type: kafka bootstrap-servers: ${KAFKA_BROKERS:114.55.58.251:9092} compression-type: zstd linger-ms: 20 diff --git a/modules/apps/xinda-push-app/src/test/java/com/lingniu/ingest/xindapushapp/XindaPushAppCompositionTest.java b/modules/apps/xinda-push-app/src/test/java/com/lingniu/ingest/xindapushapp/XindaPushAppCompositionTest.java index a68d6d05..34d7e6f1 100644 --- a/modules/apps/xinda-push-app/src/test/java/com/lingniu/ingest/xindapushapp/XindaPushAppCompositionTest.java +++ b/modules/apps/xinda-push-app/src/test/java/com/lingniu/ingest/xindapushapp/XindaPushAppCompositionTest.java @@ -9,9 +9,9 @@ import com.lingniu.ingest.inbound.xinda.config.XindaPushAutoConfiguration; import com.lingniu.ingest.sink.archive.ArchiveStore; import com.lingniu.ingest.sink.archive.RawArchiveEventSink; import com.lingniu.ingest.sink.archive.config.SinkArchiveAutoConfiguration; -import com.lingniu.ingest.sink.mq.KafkaEnvelopeDeadLetterSink; -import com.lingniu.ingest.sink.mq.KafkaEventSink; -import com.lingniu.ingest.sink.mq.SinkMqAutoConfiguration; +import com.lingniu.ingest.sink.kafka.KafkaEnvelopeDeadLetterSink; +import com.lingniu.ingest.sink.kafka.KafkaEventSink; +import com.lingniu.ingest.sink.kafka.KafkaSinkAutoConfiguration; import org.apache.kafka.clients.producer.KafkaProducer; import org.junit.jupiter.api.Test; import org.springframework.boot.autoconfigure.AutoConfigurations; @@ -26,7 +26,7 @@ class XindaPushAppCompositionTest { .withConfiguration(AutoConfigurations.of( IngestCoreAutoConfiguration.class, VehicleIdentityAutoConfiguration.class, - SinkMqAutoConfiguration.class, + KafkaSinkAutoConfiguration.class, SinkArchiveAutoConfiguration.class, XindaPushAutoConfiguration.class)) .withAllowBeanDefinitionOverriding(true) @@ -37,10 +37,11 @@ class XindaPushAppCompositionTest { "lingniu.ingest.xinda-push.port=10100", "lingniu.ingest.xinda-push.username=test", "lingniu.ingest.xinda-push.password=test", - "lingniu.ingest.sink.mq.enabled=true", - "lingniu.ingest.sink.mq.type=kafka", - "lingniu.ingest.sink.mq.bootstrap-servers=localhost:9092", - "lingniu.ingest.sink.mq.consumer.enabled=false", + "lingniu.ingest.identity.store=mysql", + "lingniu.ingest.identity.mysql.initialize-schema=false", + "lingniu.ingest.sink.kafka.enabled=true", + "lingniu.ingest.sink.kafka.bootstrap-servers=localhost:9092", + "lingniu.ingest.sink.kafka.consumer.enabled=false", "lingniu.ingest.sink.archive.enabled=true", "lingniu.ingest.sink.archive.path=target/test-archive-xinda", "lingniu.ingest.event-file-store.enabled=false", diff --git a/modules/apps/xinda-push-app/src/test/java/com/lingniu/ingest/xindapushapp/XindaPushAppDefaultsTest.java b/modules/apps/xinda-push-app/src/test/java/com/lingniu/ingest/xindapushapp/XindaPushAppDefaultsTest.java index 3c6914b1..ca28d5e5 100644 --- a/modules/apps/xinda-push-app/src/test/java/com/lingniu/ingest/xindapushapp/XindaPushAppDefaultsTest.java +++ b/modules/apps/xinda-push-app/src/test/java/com/lingniu/ingest/xindapushapp/XindaPushAppDefaultsTest.java @@ -26,13 +26,14 @@ class XindaPushAppDefaultsTest { .containsEntry("lingniu.ingest.xinda-push.port", "${XINDA_PUSH_PORT:10100}") .containsEntry("lingniu.ingest.xinda-push.username", "${XINDA_PUSH_USERNAME:}") .containsEntry("lingniu.ingest.xinda-push.subscribe-msg-ids[0]", "${XINDA_PUSH_MSG_ID_LOCATION:0200}") - .containsEntry("lingniu.ingest.sink.mq.enabled", "${KAFKA_ENABLED:true}") - .containsEntry("lingniu.ingest.sink.mq.consumer.enabled", false) - .containsEntry("lingniu.ingest.sink.mq.topics.realtime", "${KAFKA_TOPIC_XINDA_PUSH_EVENT:vehicle.event.xinda-push.v1}") - .containsEntry("lingniu.ingest.sink.mq.topics.raw-archive", "${KAFKA_TOPIC_XINDA_PUSH_RAW:vehicle.raw.xinda-push.v1}") - .containsEntry("lingniu.ingest.sink.mq.topics.dlq", "${KAFKA_TOPIC_XINDA_PUSH_DLQ:vehicle.dlq.xinda-push.v1}") - .containsEntry("lingniu.ingest.identity.store", "${VEHICLE_IDENTITY_STORE:file}") - .containsEntry("lingniu.ingest.identity.mysql.table", "${VEHICLE_IDENTITY_MYSQL_TABLE:vehicle_identity_binding}") + .containsEntry("lingniu.ingest.sink.kafka.enabled", "${KAFKA_ENABLED:true}") + .containsEntry("lingniu.ingest.sink.kafka.consumer.enabled", false) + .containsEntry("lingniu.ingest.sink.kafka.topics.realtime", "${KAFKA_TOPIC_XINDA_PUSH_EVENT:vehicle.event.xinda-push.v1}") + .containsEntry("lingniu.ingest.sink.kafka.topics.raw-archive", "${KAFKA_TOPIC_XINDA_PUSH_RAW:vehicle.raw.xinda-push.v1}") + .containsEntry("lingniu.ingest.sink.kafka.topics.dlq", "${KAFKA_TOPIC_XINDA_PUSH_DLQ:vehicle.dlq.xinda-push.v1}") + .containsEntry("lingniu.ingest.identity.store", "${VEHICLE_IDENTITY_STORE:mysql}") + .containsEntry("lingniu.ingest.identity.mysql.table-name", + "${VEHICLE_IDENTITY_MYSQL_TABLE:vehicle_identity_bindings}") .containsEntry("lingniu.ingest.sink.archive.enabled", "${SINK_ARCHIVE_ENABLED:true}") .containsEntry("lingniu.ingest.event-file-store.enabled", false) .containsEntry("lingniu.ingest.event-history.enabled", false) diff --git a/modules/apps/yutong-mqtt-app/src/main/resources/application.yml b/modules/apps/yutong-mqtt-app/src/main/resources/application.yml index f796ad74..9eefde95 100644 --- a/modules/apps/yutong-mqtt-app/src/main/resources/application.yml +++ b/modules/apps/yutong-mqtt-app/src/main/resources/application.yml @@ -69,9 +69,8 @@ lingniu: table-name: ${VEHICLE_IDENTITY_MYSQL_TABLE:vehicle_identity_bindings} initialize-schema: ${VEHICLE_IDENTITY_MYSQL_INITIALIZE_SCHEMA:true} sink: - mq: + kafka: enabled: ${KAFKA_ENABLED:true} - type: kafka bootstrap-servers: ${KAFKA_BROKERS:114.55.58.251:9092} compression-type: zstd linger-ms: 20 diff --git a/modules/apps/yutong-mqtt-app/src/test/java/com/lingniu/ingest/yutongmqttapp/YutongMqttAppCompositionTest.java b/modules/apps/yutong-mqtt-app/src/test/java/com/lingniu/ingest/yutongmqttapp/YutongMqttAppCompositionTest.java index 33b8aba0..73ae6312 100644 --- a/modules/apps/yutong-mqtt-app/src/test/java/com/lingniu/ingest/yutongmqttapp/YutongMqttAppCompositionTest.java +++ b/modules/apps/yutong-mqtt-app/src/test/java/com/lingniu/ingest/yutongmqttapp/YutongMqttAppCompositionTest.java @@ -10,9 +10,9 @@ import com.lingniu.ingest.inbound.mqtt.profile.MqttProfileRegistry; import com.lingniu.ingest.sink.archive.ArchiveStore; import com.lingniu.ingest.sink.archive.RawArchiveEventSink; import com.lingniu.ingest.sink.archive.config.SinkArchiveAutoConfiguration; -import com.lingniu.ingest.sink.mq.KafkaEnvelopeDeadLetterSink; -import com.lingniu.ingest.sink.mq.KafkaEventSink; -import com.lingniu.ingest.sink.mq.SinkMqAutoConfiguration; +import com.lingniu.ingest.sink.kafka.KafkaEnvelopeDeadLetterSink; +import com.lingniu.ingest.sink.kafka.KafkaEventSink; +import com.lingniu.ingest.sink.kafka.KafkaSinkAutoConfiguration; import org.apache.kafka.clients.producer.KafkaProducer; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; @@ -31,7 +31,7 @@ class YutongMqttAppCompositionTest { .withConfiguration(AutoConfigurations.of( IngestCoreAutoConfiguration.class, VehicleIdentityAutoConfiguration.class, - SinkMqAutoConfiguration.class, + KafkaSinkAutoConfiguration.class, SinkArchiveAutoConfiguration.class, MqttInboundAutoConfiguration.class)) .withAllowBeanDefinitionOverriding(true) @@ -45,10 +45,9 @@ class YutongMqttAppCompositionTest { "lingniu.ingest.mqtt.endpoints[0].profile=yutong", "lingniu.ingest.identity.store=mysql", "lingniu.ingest.identity.mysql.initialize-schema=false", - "lingniu.ingest.sink.mq.enabled=true", - "lingniu.ingest.sink.mq.type=kafka", - "lingniu.ingest.sink.mq.bootstrap-servers=localhost:9092", - "lingniu.ingest.sink.mq.consumer.enabled=false", + "lingniu.ingest.sink.kafka.enabled=true", + "lingniu.ingest.sink.kafka.bootstrap-servers=localhost:9092", + "lingniu.ingest.sink.kafka.consumer.enabled=false", "lingniu.ingest.sink.archive.enabled=true", "lingniu.ingest.sink.archive.path=target/test-archive-yutong", "lingniu.ingest.event-history.enabled=false"); diff --git a/modules/apps/yutong-mqtt-app/src/test/java/com/lingniu/ingest/yutongmqttapp/YutongMqttAppDefaultsTest.java b/modules/apps/yutong-mqtt-app/src/test/java/com/lingniu/ingest/yutongmqttapp/YutongMqttAppDefaultsTest.java index 58e0128f..f5358a0f 100644 --- a/modules/apps/yutong-mqtt-app/src/test/java/com/lingniu/ingest/yutongmqttapp/YutongMqttAppDefaultsTest.java +++ b/modules/apps/yutong-mqtt-app/src/test/java/com/lingniu/ingest/yutongmqttapp/YutongMqttAppDefaultsTest.java @@ -28,11 +28,11 @@ class YutongMqttAppDefaultsTest { .containsEntry("lingniu.ingest.mqtt.endpoints[0].uri", "${YUTONG_MQTT_URI:}") .containsEntry("lingniu.ingest.mqtt.endpoints[0].topic", "${YUTONG_MQTT_TOPIC:#}") .containsEntry("lingniu.ingest.mqtt.endpoints[0].profile", "yutong") - .containsEntry("lingniu.ingest.sink.mq.enabled", "${KAFKA_ENABLED:true}") - .containsEntry("lingniu.ingest.sink.mq.consumer.enabled", false) - .containsEntry("lingniu.ingest.sink.mq.topics.realtime", "${KAFKA_TOPIC_YUTONG_MQTT_EVENT:vehicle.event.mqtt-yutong.v1}") - .containsEntry("lingniu.ingest.sink.mq.topics.raw-archive", "${KAFKA_TOPIC_YUTONG_MQTT_RAW:vehicle.raw.mqtt-yutong.v1}") - .containsEntry("lingniu.ingest.sink.mq.topics.dlq", "${KAFKA_TOPIC_YUTONG_MQTT_DLQ:vehicle.dlq.mqtt-yutong.v1}") + .containsEntry("lingniu.ingest.sink.kafka.enabled", "${KAFKA_ENABLED:true}") + .containsEntry("lingniu.ingest.sink.kafka.consumer.enabled", false) + .containsEntry("lingniu.ingest.sink.kafka.topics.realtime", "${KAFKA_TOPIC_YUTONG_MQTT_EVENT:vehicle.event.mqtt-yutong.v1}") + .containsEntry("lingniu.ingest.sink.kafka.topics.raw-archive", "${KAFKA_TOPIC_YUTONG_MQTT_RAW:vehicle.raw.mqtt-yutong.v1}") + .containsEntry("lingniu.ingest.sink.kafka.topics.dlq", "${KAFKA_TOPIC_YUTONG_MQTT_DLQ:vehicle.dlq.mqtt-yutong.v1}") .containsEntry("lingniu.ingest.identity.store", "${VEHICLE_IDENTITY_STORE:mysql}") .containsEntry("lingniu.ingest.identity.mysql.table-name", "${VEHICLE_IDENTITY_MYSQL_TABLE:vehicle_identity_bindings}") diff --git a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestor.java b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestor.java index 04bf9a1f..00ac451f 100644 --- a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestor.java +++ b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestor.java @@ -5,7 +5,7 @@ import com.lingniu.ingest.api.consumer.EnvelopeBatchIngestor; import com.lingniu.ingest.api.consumer.EnvelopeIngestResult; import com.lingniu.ingest.api.history.EventFileRecord; import com.lingniu.ingest.api.history.EventFileStore; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import com.lingniu.ingest.tdenginehistory.TdengineEnvelopeRows; import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter; import org.slf4j.Logger; diff --git a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/TelemetryEnvelopeRecordMapper.java b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/TelemetryEnvelopeRecordMapper.java index ffef1246..2c22d708 100644 --- a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/TelemetryEnvelopeRecordMapper.java +++ b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/TelemetryEnvelopeRecordMapper.java @@ -5,8 +5,8 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.lingniu.ingest.api.ProtocolId; import com.lingniu.ingest.api.event.RawArchiveKeys; import com.lingniu.ingest.api.history.EventFileRecord; -import com.lingniu.ingest.sink.mq.proto.RawArchiveRef; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.RawArchiveRef; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import com.google.protobuf.InvalidProtocolBufferException; import com.google.protobuf.util.JsonFormat; diff --git a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfiguration.java b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfiguration.java index e1f444ef..47b5df20 100644 --- a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfiguration.java +++ b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfiguration.java @@ -15,7 +15,7 @@ import com.lingniu.ingest.eventhistory.TelemetryFieldHistoryController; import com.lingniu.ingest.eventhistory.TelemetryEnvelopeRecordMapper; import com.lingniu.ingest.protocol.gb32960.codec.Gb32960MessageDecoder; import com.lingniu.ingest.protocol.gb32960.config.Gb32960AutoConfiguration; -import com.lingniu.ingest.sink.mq.SinkMqAutoConfiguration; +import com.lingniu.ingest.sink.kafka.KafkaSinkAutoConfiguration; import com.lingniu.ingest.tdenginehistory.config.TdengineHistoryAutoConfiguration; import com.lingniu.ingest.tdenginehistory.TdengineHistoryReader; import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter; @@ -41,12 +41,12 @@ import java.nio.file.Path; * * *

当前生产 history app 以 TDengine 为准;Kafka consumer 是否启动还取决于 - * {@code lingniu.ingest.sink.mq.consumer.enabled=true}。 + * {@code lingniu.ingest.sink.kafka.consumer.enabled=true}。 */ @AutoConfiguration @AutoConfigureAfter({ Gb32960AutoConfiguration.class, - SinkMqAutoConfiguration.class, + KafkaSinkAutoConfiguration.class, TdengineHistoryAutoConfiguration.class }) @ConditionalOnProperty(prefix = "lingniu.ingest.event-history", name = "enabled", havingValue = "true") @@ -82,7 +82,7 @@ public class EventHistoryAutoConfiguration { @ConditionalOnMissingBean(name = "eventHistoryEnvelopeConsumerProcessor") public EnvelopeConsumerProcessor eventHistoryEnvelopeConsumerProcessor(EventHistoryEnvelopeIngestor ingestor, EnvelopeDeadLetterSink deadLetterSink) { - // 只注册 processor;真正拉 Kafka 的 runner 在 sink-mq 模块按 consumer.enabled 决定是否创建。 + // 只注册 processor;真正拉 Kafka 的 runner 按 consumer.enabled 决定是否创建。 return new EnvelopeConsumerProcessor("event-history", ingestor, deadLetterSink); } diff --git a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestorTest.java b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestorTest.java index fd89b197..084f688b 100644 --- a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestorTest.java +++ b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestorTest.java @@ -9,13 +9,13 @@ import com.lingniu.ingest.api.event.RawArchiveKeys; import com.lingniu.ingest.api.history.EventFileQuery; import com.lingniu.ingest.api.history.EventFileRecord; import com.lingniu.ingest.api.history.EventFileStore; -import com.lingniu.ingest.sink.mq.proto.RawArchiveRef; -import com.lingniu.ingest.sink.mq.proto.RawFrameFactPayload; -import com.lingniu.ingest.sink.mq.proto.TelemetryField; -import com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; -import com.lingniu.ingest.sink.mq.proto.LocationPayload; -import com.lingniu.ingest.sink.mq.proto.ParseStatusProto; +import com.lingniu.ingest.sink.kafka.proto.RawArchiveRef; +import com.lingniu.ingest.sink.kafka.proto.RawFrameFactPayload; +import com.lingniu.ingest.sink.kafka.proto.TelemetryField; +import com.lingniu.ingest.sink.kafka.proto.TelemetrySnapshot; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.LocationPayload; +import com.lingniu.ingest.sink.kafka.proto.ParseStatusProto; import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter; import com.lingniu.ingest.tdenginehistory.TdengineLocationRow; import com.lingniu.ingest.tdenginehistory.TdengineRawFrameRow; diff --git a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/TelemetryEnvelopeRecordMapperTest.java b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/TelemetryEnvelopeRecordMapperTest.java index 14f6d38e..96538c3a 100644 --- a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/TelemetryEnvelopeRecordMapperTest.java +++ b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/TelemetryEnvelopeRecordMapperTest.java @@ -2,10 +2,10 @@ package com.lingniu.ingest.eventhistory; import com.lingniu.ingest.api.ProtocolId; import com.lingniu.ingest.api.history.EventFileRecord; -import com.lingniu.ingest.sink.mq.proto.RawArchiveRef; -import com.lingniu.ingest.sink.mq.proto.TelemetryField; -import com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.RawArchiveRef; +import com.lingniu.ingest.sink.kafka.proto.TelemetryField; +import com.lingniu.ingest.sink.kafka.proto.TelemetrySnapshot; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import com.fasterxml.jackson.databind.ObjectMapper; import org.junit.jupiter.api.Test; diff --git a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfigurationTest.java b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfigurationTest.java index c855b934..40b80ca3 100644 --- a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfigurationTest.java +++ b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfigurationTest.java @@ -37,7 +37,7 @@ class EventHistoryAutoConfigurationTest { contextRunner .withPropertyValues( "lingniu.ingest.event-history.enabled=true", - "lingniu.ingest.sink.mq.consumer.enabled=true") + "lingniu.ingest.sink.kafka.consumer.enabled=true") .run(context -> { assertThat(context).hasSingleBean(TelemetryEnvelopeRecordMapper.class); assertThat(context).hasSingleBean(EventHistoryEnvelopeIngestor.class); diff --git a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/config/EventHistoryPomBoundaryTest.java b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/config/EventHistoryPomBoundaryTest.java index 095d2248..388f073a 100644 --- a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/config/EventHistoryPomBoundaryTest.java +++ b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/config/EventHistoryPomBoundaryTest.java @@ -22,7 +22,7 @@ class EventHistoryPomBoundaryTest { } @Test - void kafkaClientsIsOwnedBySinkMqModule() throws Exception { + void kafkaClientsIsOwnedByKafkaSinkModule() throws Exception { for (Path pom : servicePoms()) { assertThat(directDependencies(readPom(pom))) .as(pom.toString()) diff --git a/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/VehicleStatEnvelopeIngestor.java b/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/VehicleStatEnvelopeIngestor.java index fa00896d..35936a5c 100644 --- a/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/VehicleStatEnvelopeIngestor.java +++ b/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/VehicleStatEnvelopeIngestor.java @@ -3,7 +3,7 @@ package com.lingniu.ingest.vehiclestat; import com.google.protobuf.InvalidProtocolBufferException; import com.lingniu.ingest.api.consumer.EnvelopeIngestResult; import com.lingniu.ingest.api.consumer.EnvelopeIngestor; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import com.lingniu.ingest.vehiclestat.jt808.Jt808MileageStreamProcessor; public final class VehicleStatEnvelopeIngestor implements EnvelopeIngestor { diff --git a/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/config/VehicleStatAutoConfiguration.java b/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/config/VehicleStatAutoConfiguration.java index 4bfdffa6..8e1c069a 100644 --- a/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/config/VehicleStatAutoConfiguration.java +++ b/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/config/VehicleStatAutoConfiguration.java @@ -78,7 +78,7 @@ public class VehicleStatAutoConfiguration { @ConditionalOnMissingBean(name = "vehicleStatEnvelopeConsumerProcessor") public EnvelopeConsumerProcessor vehicleStatEnvelopeConsumerProcessor(VehicleStatEnvelopeIngestor ingestor, EnvelopeDeadLetterSink deadLetterSink) { - // Bean 名必须和 sink-mq 默认 binding 对齐,KafkaEnvelopeConsumerFactory 才能自动创建 worker。 + // Bean 名必须和 Kafka 默认 binding 对齐,KafkaEnvelopeConsumerFactory 才能自动创建 worker。 return new EnvelopeConsumerProcessor("vehicle-stat", ingestor, deadLetterSink); } } diff --git a/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/jt808/Jt808LocationPointExtractor.java b/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/jt808/Jt808LocationPointExtractor.java index cd553d68..8883391d 100644 --- a/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/jt808/Jt808LocationPointExtractor.java +++ b/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/jt808/Jt808LocationPointExtractor.java @@ -1,7 +1,7 @@ package com.lingniu.ingest.vehiclestat.jt808; -import com.lingniu.ingest.sink.mq.proto.TelemetryField; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.TelemetryField; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import java.time.Instant; import java.util.LinkedHashMap; diff --git a/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/jt808/Jt808MileageStreamProcessor.java b/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/jt808/Jt808MileageStreamProcessor.java index d4683a1c..05fba77e 100644 --- a/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/jt808/Jt808MileageStreamProcessor.java +++ b/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/jt808/Jt808MileageStreamProcessor.java @@ -1,6 +1,6 @@ package com.lingniu.ingest.vehiclestat.jt808; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import com.lingniu.ingest.vehiclestat.VehicleStatRepository; import java.time.LocalDate; diff --git a/modules/services/vehicle-stat-service/src/test/java/com/lingniu/ingest/vehiclestat/VehicleStatEnvelopeIngestorTest.java b/modules/services/vehicle-stat-service/src/test/java/com/lingniu/ingest/vehiclestat/VehicleStatEnvelopeIngestorTest.java index ea96483c..e13d6051 100644 --- a/modules/services/vehicle-stat-service/src/test/java/com/lingniu/ingest/vehiclestat/VehicleStatEnvelopeIngestorTest.java +++ b/modules/services/vehicle-stat-service/src/test/java/com/lingniu/ingest/vehiclestat/VehicleStatEnvelopeIngestorTest.java @@ -2,10 +2,10 @@ package com.lingniu.ingest.vehiclestat; import com.lingniu.ingest.api.consumer.EnvelopeIngestResult; import com.lingniu.ingest.api.consumer.EnvelopeIngestor; -import com.lingniu.ingest.sink.mq.proto.RawArchiveRef; -import com.lingniu.ingest.sink.mq.proto.TelemetryField; -import com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.RawArchiveRef; +import com.lingniu.ingest.sink.kafka.proto.TelemetryField; +import com.lingniu.ingest.sink.kafka.proto.TelemetrySnapshot; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import com.lingniu.ingest.vehiclestat.jt808.Jt808MileageStreamProcessor; import org.junit.jupiter.api.Test; diff --git a/modules/services/vehicle-stat-service/src/test/java/com/lingniu/ingest/vehiclestat/jt808/Jt808LocationPointExtractorTest.java b/modules/services/vehicle-stat-service/src/test/java/com/lingniu/ingest/vehiclestat/jt808/Jt808LocationPointExtractorTest.java index a6d555d7..579c7135 100644 --- a/modules/services/vehicle-stat-service/src/test/java/com/lingniu/ingest/vehiclestat/jt808/Jt808LocationPointExtractorTest.java +++ b/modules/services/vehicle-stat-service/src/test/java/com/lingniu/ingest/vehiclestat/jt808/Jt808LocationPointExtractorTest.java @@ -1,8 +1,8 @@ package com.lingniu.ingest.vehiclestat.jt808; -import com.lingniu.ingest.sink.mq.proto.TelemetryField; -import com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.TelemetryField; +import com.lingniu.ingest.sink.kafka.proto.TelemetrySnapshot; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import org.junit.jupiter.api.Test; import java.time.Instant; diff --git a/modules/services/vehicle-stat-service/src/test/java/com/lingniu/ingest/vehiclestat/jt808/Jt808MileageStreamProcessorTest.java b/modules/services/vehicle-stat-service/src/test/java/com/lingniu/ingest/vehiclestat/jt808/Jt808MileageStreamProcessorTest.java index fe972e7d..6ef612e5 100644 --- a/modules/services/vehicle-stat-service/src/test/java/com/lingniu/ingest/vehiclestat/jt808/Jt808MileageStreamProcessorTest.java +++ b/modules/services/vehicle-stat-service/src/test/java/com/lingniu/ingest/vehiclestat/jt808/Jt808MileageStreamProcessorTest.java @@ -1,8 +1,8 @@ package com.lingniu.ingest.vehiclestat.jt808; -import com.lingniu.ingest.sink.mq.proto.TelemetryField; -import com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.TelemetryField; +import com.lingniu.ingest.sink.kafka.proto.TelemetrySnapshot; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import com.lingniu.ingest.vehiclestat.DailyMileageStrategy; import com.lingniu.ingest.vehiclestat.VehicleDailyStatResult; import com.lingniu.ingest.vehiclestat.VehicleStatRepository; diff --git a/modules/services/vehicle-state-service/src/main/java/com/lingniu/ingest/vehiclestate/VehicleStateEnvelopeIngestor.java b/modules/services/vehicle-state-service/src/main/java/com/lingniu/ingest/vehiclestate/VehicleStateEnvelopeIngestor.java index f4e4489e..0411a863 100644 --- a/modules/services/vehicle-state-service/src/main/java/com/lingniu/ingest/vehiclestate/VehicleStateEnvelopeIngestor.java +++ b/modules/services/vehicle-state-service/src/main/java/com/lingniu/ingest/vehiclestate/VehicleStateEnvelopeIngestor.java @@ -3,7 +3,7 @@ package com.lingniu.ingest.vehiclestate; import com.google.protobuf.InvalidProtocolBufferException; import com.lingniu.ingest.api.consumer.EnvelopeIngestResult; import com.lingniu.ingest.api.consumer.EnvelopeIngestor; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; public final class VehicleStateEnvelopeIngestor implements EnvelopeIngestor { diff --git a/modules/services/vehicle-state-service/src/main/java/com/lingniu/ingest/vehiclestate/VehicleStateUpdater.java b/modules/services/vehicle-state-service/src/main/java/com/lingniu/ingest/vehiclestate/VehicleStateUpdater.java index 699e36c8..8d6d6788 100644 --- a/modules/services/vehicle-state-service/src/main/java/com/lingniu/ingest/vehiclestate/VehicleStateUpdater.java +++ b/modules/services/vehicle-state-service/src/main/java/com/lingniu/ingest/vehiclestate/VehicleStateUpdater.java @@ -1,7 +1,7 @@ package com.lingniu.ingest.vehiclestate; -import com.lingniu.ingest.sink.mq.proto.TelemetryField; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.TelemetryField; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import java.util.LinkedHashMap; import java.util.Map; diff --git a/modules/services/vehicle-state-service/src/main/java/com/lingniu/ingest/vehiclestate/config/VehicleStateAutoConfiguration.java b/modules/services/vehicle-state-service/src/main/java/com/lingniu/ingest/vehiclestate/config/VehicleStateAutoConfiguration.java index d65dd316..c808e737 100644 --- a/modules/services/vehicle-state-service/src/main/java/com/lingniu/ingest/vehiclestate/config/VehicleStateAutoConfiguration.java +++ b/modules/services/vehicle-state-service/src/main/java/com/lingniu/ingest/vehiclestate/config/VehicleStateAutoConfiguration.java @@ -45,7 +45,7 @@ public class VehicleStateAutoConfiguration { @ConditionalOnMissingBean(name = "vehicleStateEnvelopeConsumerProcessor") public EnvelopeConsumerProcessor vehicleStateEnvelopeConsumerProcessor(VehicleStateEnvelopeIngestor ingestor, EnvelopeDeadLetterSink deadLetterSink) { - // Bean 名必须和 sink-mq 默认 binding 对齐,KafkaEnvelopeConsumerFactory 才能自动创建 worker。 + // Bean 名必须和 Kafka 默认 binding 对齐,KafkaEnvelopeConsumerFactory 才能自动创建 worker。 return new EnvelopeConsumerProcessor("vehicle-state", ingestor, deadLetterSink); } diff --git a/modules/services/vehicle-state-service/src/test/java/com/lingniu/ingest/vehiclestate/VehicleStateEnvelopeIngestorTest.java b/modules/services/vehicle-state-service/src/test/java/com/lingniu/ingest/vehiclestate/VehicleStateEnvelopeIngestorTest.java index 6f5ec96b..dbf9dd92 100644 --- a/modules/services/vehicle-state-service/src/test/java/com/lingniu/ingest/vehiclestate/VehicleStateEnvelopeIngestorTest.java +++ b/modules/services/vehicle-state-service/src/test/java/com/lingniu/ingest/vehiclestate/VehicleStateEnvelopeIngestorTest.java @@ -2,10 +2,10 @@ package com.lingniu.ingest.vehiclestate; import com.lingniu.ingest.api.consumer.EnvelopeIngestResult; import com.lingniu.ingest.api.consumer.EnvelopeIngestor; -import com.lingniu.ingest.sink.mq.proto.RawArchiveRef; -import com.lingniu.ingest.sink.mq.proto.TelemetryField; -import com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.RawArchiveRef; +import com.lingniu.ingest.sink.kafka.proto.TelemetryField; +import com.lingniu.ingest.sink.kafka.proto.TelemetrySnapshot; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import org.junit.jupiter.api.Test; import java.util.Optional; diff --git a/modules/services/vehicle-state-service/src/test/java/com/lingniu/ingest/vehiclestate/VehicleStateUpdaterTest.java b/modules/services/vehicle-state-service/src/test/java/com/lingniu/ingest/vehiclestate/VehicleStateUpdaterTest.java index 7f93be7c..c896c85b 100644 --- a/modules/services/vehicle-state-service/src/test/java/com/lingniu/ingest/vehiclestate/VehicleStateUpdaterTest.java +++ b/modules/services/vehicle-state-service/src/test/java/com/lingniu/ingest/vehiclestate/VehicleStateUpdaterTest.java @@ -1,8 +1,8 @@ package com.lingniu.ingest.vehiclestate; -import com.lingniu.ingest.sink.mq.proto.TelemetryField; -import com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.TelemetryField; +import com.lingniu.ingest.sink.kafka.proto.TelemetrySnapshot; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import org.junit.jupiter.api.Test; import static org.assertj.core.api.Assertions.assertThat; diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/EnvelopeMapper.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/EnvelopeMapper.java similarity index 86% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/EnvelopeMapper.java rename to modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/EnvelopeMapper.java index 2398d5e4..405da958 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/EnvelopeMapper.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/EnvelopeMapper.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import com.lingniu.ingest.api.event.AlarmPayload; import com.lingniu.ingest.api.event.LocationPayload; @@ -10,10 +10,10 @@ import com.lingniu.ingest.api.event.VehicleEvent; import com.lingniu.ingest.api.event.VehicleEventTelemetrySnapshotMapper; import com.lingniu.ingest.facts.FactIds; import com.lingniu.ingest.facts.VehicleKey; -import com.lingniu.ingest.sink.mq.proto.ParseStatusProto; -import com.lingniu.ingest.sink.mq.proto.TelemetryField; -import com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot.Builder; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.ParseStatusProto; +import com.lingniu.ingest.sink.kafka.proto.TelemetryField; +import com.lingniu.ingest.sink.kafka.proto.TelemetrySnapshot.Builder; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import java.security.MessageDigest; import java.security.NoSuchAlgorithmException; @@ -53,23 +53,23 @@ public final class EnvelopeMapper { case VehicleEvent.Location l -> b.setLocation(buildLocation(l.payload())); case VehicleEvent.Alarm a -> b.setAlarm(buildAlarm(a.payload())); case VehicleEvent.Login lg -> b.setLogin( - com.lingniu.ingest.sink.mq.proto.LoginPayload.newBuilder() + com.lingniu.ingest.sink.kafka.proto.LoginPayload.newBuilder() .setIccid(nullToEmpty(lg.iccid())) .setProtocolVersion(nullToEmpty(lg.protocolVersion())) .build()); case VehicleEvent.Logout ignored -> b.setLogout( - com.lingniu.ingest.sink.mq.proto.LogoutPayload.getDefaultInstance()); + com.lingniu.ingest.sink.kafka.proto.LogoutPayload.getDefaultInstance()); case VehicleEvent.Heartbeat ignored -> b.setHeartbeat( - com.lingniu.ingest.sink.mq.proto.HeartbeatPayload.getDefaultInstance()); + com.lingniu.ingest.sink.kafka.proto.HeartbeatPayload.getDefaultInstance()); case VehicleEvent.MediaMeta m -> b.setMediaMeta( - com.lingniu.ingest.sink.mq.proto.MediaMetaPayload.newBuilder() + com.lingniu.ingest.sink.kafka.proto.MediaMetaPayload.newBuilder() .setMediaId(nullToEmpty(m.mediaId())) .setMediaType(nullToEmpty(m.mediaType())) .setSizeBytes(m.sizeBytes()) .setArchiveRef(nullToEmpty(m.archiveRef())) .build()); case VehicleEvent.Passthrough p -> b.setPassthrough( - com.lingniu.ingest.sink.mq.proto.PassthroughPayload.newBuilder() + com.lingniu.ingest.sink.kafka.proto.PassthroughPayload.newBuilder() .setPassthroughType(p.passthroughType()) .setData(com.google.protobuf.ByteString.copyFrom( p.data() == null ? new byte[0] : p.data())) @@ -88,13 +88,13 @@ public final class EnvelopeMapper { b.putMetadata(RawArchiveKeys.META_KEY, key); b.putMetadata(RawArchiveKeys.META_URI, uri); b.putMetadata(RawArchiveKeys.META_EVENT_ID, ra.eventId()); - b.setRawArchive(com.lingniu.ingest.sink.mq.proto.RawArchiveRef.newBuilder() + b.setRawArchive(com.lingniu.ingest.sink.kafka.proto.RawArchiveRef.newBuilder() .setUri(uri) .setChecksum(checksum) .setSizeBytes(size) .setParsedJson(nullToEmpty(ra.parsedJson())) .build()); - b.setRawFrameFact(com.lingniu.ingest.sink.mq.proto.RawFrameFactPayload.newBuilder() + b.setRawFrameFact(com.lingniu.ingest.sink.kafka.proto.RawFrameFactPayload.newBuilder() .setFrameId(frameId) .setVehicleKey(vehicleKey) .setVin(nullToEmpty(ra.vin())) @@ -114,8 +114,8 @@ public final class EnvelopeMapper { return b.build(); } - private static com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot buildTelemetrySnapshot(TelemetrySnapshot snapshot) { - Builder b = com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot.newBuilder() + private static com.lingniu.ingest.sink.kafka.proto.TelemetrySnapshot buildTelemetrySnapshot(TelemetrySnapshot snapshot) { + Builder b = com.lingniu.ingest.sink.kafka.proto.TelemetrySnapshot.newBuilder() .setEventType(snapshot.eventType()) .setRawArchiveUri(snapshot.rawArchiveUri()); for (TelemetryFieldValue field : snapshot.fields()) { @@ -131,8 +131,8 @@ public final class EnvelopeMapper { return b.build(); } - private static com.lingniu.ingest.sink.mq.proto.RealtimePayload buildRealtime(RealtimePayload p) { - var b = com.lingniu.ingest.sink.mq.proto.RealtimePayload.newBuilder(); + private static com.lingniu.ingest.sink.kafka.proto.RealtimePayload buildRealtime(RealtimePayload p) { + var b = com.lingniu.ingest.sink.kafka.proto.RealtimePayload.newBuilder(); if (p.speedKmh() != null) b.setSpeedKmh(p.speedKmh()); if (p.totalMileageKm() != null) b.setTotalMileageKm(p.totalMileageKm()); if (p.batterySoc() != null) b.setBatterySoc(p.batterySoc()); @@ -159,8 +159,8 @@ public final class EnvelopeMapper { return b.build(); } - private static com.lingniu.ingest.sink.mq.proto.LocationPayload buildLocation(LocationPayload p) { - return com.lingniu.ingest.sink.mq.proto.LocationPayload.newBuilder() + private static com.lingniu.ingest.sink.kafka.proto.LocationPayload buildLocation(LocationPayload p) { + return com.lingniu.ingest.sink.kafka.proto.LocationPayload.newBuilder() .setLongitude(p.longitude()) .setLatitude(p.latitude()) .setAltitudeM(p.altitudeM()) @@ -171,8 +171,8 @@ public final class EnvelopeMapper { .build(); } - private static com.lingniu.ingest.sink.mq.proto.AlarmPayload buildAlarm(AlarmPayload p) { - var b = com.lingniu.ingest.sink.mq.proto.AlarmPayload.newBuilder() + private static com.lingniu.ingest.sink.kafka.proto.AlarmPayload buildAlarm(AlarmPayload p) { + var b = com.lingniu.ingest.sink.kafka.proto.AlarmPayload.newBuilder() .setLevel(p.level().name()) .setAlarmTypeCode(p.alarmTypeCode()) .setAlarmTypeName(nullToEmpty(p.alarmTypeName())); diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerFactory.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerFactory.java similarity index 83% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerFactory.java rename to modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerFactory.java index 7fda1bc0..ed83e261 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerFactory.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerFactory.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor; import org.apache.kafka.clients.consumer.ConsumerConfig; @@ -28,15 +28,15 @@ public final class KafkaEnvelopeConsumerFactory { } public List createWorkers(Map processors, - SinkMqProperties props) { - Map bindings = effectiveBindings(props); + KafkaSinkProperties props) { + Map bindings = effectiveBindings(props); List workers = new ArrayList<>(); - for (Map.Entry entry : bindings.entrySet()) { + for (Map.Entry entry : bindings.entrySet()) { String processorBeanName = entry.getKey(); // binding 的 key 必须和 Spring Bean 名一致;这样配置只声明 topic/group, // 实际处理逻辑仍由各业务模块自己的 EnvelopeConsumerProcessor 承接。 EnvelopeConsumerProcessor processor = processors.get(processorBeanName); - SinkMqProperties.Binding binding = entry.getValue(); + KafkaSinkProperties.Binding binding = entry.getValue(); if (processor == null || binding == null || !binding.isEnabled()) { continue; } @@ -57,15 +57,15 @@ public final class KafkaEnvelopeConsumerFactory { return workers; } - private Map effectiveBindings(SinkMqProperties props) { - Map configured = props.getConsumer().getBindings(); + private Map effectiveBindings(KafkaSinkProperties props) { + Map configured = props.getConsumer().getBindings(); if (configured != null && !configured.isEmpty()) { return configured; } // 默认绑定仅给未显式配置 bindings 的轻量运行时兜底。 // 生产 history app 会显式绑定各协议 event/raw topic,并写入 TDengine raw_frames/locations。 - SinkMqProperties.Topics topics = props.getTopics(); - Map defaults = new LinkedHashMap<>(); + KafkaSinkProperties.Topics topics = props.getTopics(); + Map defaults = new LinkedHashMap<>(); defaults.put("eventHistoryEnvelopeConsumerProcessor", binding( "vehicle-event-history", topics.getRealtime(), topics.getLocation(), topics.getAlarm(), topics.getSession(), topics.getMediaMeta())); @@ -78,8 +78,8 @@ public final class KafkaEnvelopeConsumerFactory { return defaults; } - private SinkMqProperties.Binding binding(String groupId, String... topics) { - SinkMqProperties.Binding binding = new SinkMqProperties.Binding(); + private KafkaSinkProperties.Binding binding(String groupId, String... topics) { + KafkaSinkProperties.Binding binding = new KafkaSinkProperties.Binding(); binding.setGroupId(groupId); binding.setTopics(List.of(topics)); return binding; @@ -106,12 +106,12 @@ public final class KafkaEnvelopeConsumerFactory { return List.copyOf(clean); } - private Properties consumerProperties(SinkMqProperties props, - SinkMqProperties.Binding binding, + private Properties consumerProperties(KafkaSinkProperties props, + KafkaSinkProperties.Binding binding, String processorBeanName, int workerIndex, int concurrency) { - SinkMqProperties.Consumer consumer = props.getConsumer(); + KafkaSinkProperties.Consumer consumer = props.getConsumer(); Properties p = new Properties(); p.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, props.getBootstrapServers()); p.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); @@ -133,7 +133,7 @@ public final class KafkaEnvelopeConsumerFactory { return concurrency <= 1 ? base : base + "-" + workerIndex; } - private String groupId(SinkMqProperties.Binding binding, String processorBeanName) { + private String groupId(KafkaSinkProperties.Binding binding, String processorBeanName) { if (binding.getGroupId() != null && !binding.getGroupId().isBlank()) { return binding.getGroupId(); } diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerRunner.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerRunner.java similarity index 99% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerRunner.java rename to modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerRunner.java index 3c68b380..17e90726 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerRunner.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerRunner.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerWorker.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerWorker.java similarity index 98% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerWorker.java rename to modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerWorker.java index d640afca..b1f9e9e3 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerWorker.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerWorker.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor; import com.lingniu.ingest.api.consumer.EnvelopeConsumerRecord; diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeDeadLetterSink.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeDeadLetterSink.java similarity index 98% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeDeadLetterSink.java rename to modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeDeadLetterSink.java index e3043f68..4ba2bf2c 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeDeadLetterSink.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeDeadLetterSink.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import com.lingniu.ingest.api.consumer.EnvelopeDeadLetterRecord; import com.lingniu.ingest.api.consumer.EnvelopeDeadLetterSink; diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEventSink.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEventSink.java similarity index 99% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEventSink.java rename to modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEventSink.java index 5026004f..cf879fce 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEventSink.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaEventSink.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import com.lingniu.ingest.api.event.VehicleEvent; import com.lingniu.ingest.api.sink.EventSink; diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqAutoConfiguration.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkAutoConfiguration.java similarity index 64% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqAutoConfiguration.java rename to modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkAutoConfiguration.java index b89feaac..af7ff127 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqAutoConfiguration.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkAutoConfiguration.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import io.github.resilience4j.circuitbreaker.CircuitBreaker; import io.github.resilience4j.circuitbreaker.CircuitBreakerConfig; @@ -16,40 +16,31 @@ import java.time.Duration; import java.util.Properties; /** - * MQ Sink 自动装配。 + * Kafka Sink 自动装配。 * - *

装配前提(两个条件都要满足): - *

    - *
  1. {@code lingniu.ingest.sink.mq.enabled=true}(默认 true)—— 总开关, - * 设为 false 时本模块完全不装配,ingest-core 的 DisruptorEventBus 仍然运行但 - * 没有 Kafka sink。 - *
  2. {@code lingniu.ingest.sink.mq.type=kafka}(默认 kafka)—— 唯一生产 MQ 后端。 - *
- * - *

Producer 和 Consumer 是两个独立开关:{@code sink.mq.enabled=true} 只表示可以创建 - * Kafka producer/sink;是否从 Kafka 拉 envelope 还要单独开启 - * {@code lingniu.ingest.sink.mq.consumer.enabled=true} 并配置 bindings。 + *

{@code lingniu.ingest.sink.kafka.enabled=false} 时本模块完全不装配;Producer 和 Consumer + * 是两个独立开关,消费端还需要开启 {@code lingniu.ingest.sink.kafka.consumer.enabled=true} + * 并配置 bindings。 */ @AutoConfiguration -@EnableConfigurationProperties(SinkMqProperties.class) -@ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq", name = "enabled", havingValue = "true", matchIfMissing = true) -public class SinkMqAutoConfiguration { +@EnableConfigurationProperties(KafkaSinkProperties.class) +@ConditionalOnProperty(prefix = "lingniu.ingest.sink.kafka", name = "enabled", havingValue = "true", matchIfMissing = true) +public class KafkaSinkAutoConfiguration { @Bean @ConditionalOnMissingBean - public EnvelopeMapper envelopeMapper(SinkMqProperties props) { + public EnvelopeMapper envelopeMapper(KafkaSinkProperties props) { return new EnvelopeMapper(props.getNodeId()); } @Bean @ConditionalOnMissingBean - public TopicRouter topicRouter(SinkMqProperties props) { + public TopicRouter topicRouter(KafkaSinkProperties props) { return new TopicRouter(props.getTopics()); } @Bean @ConditionalOnMissingBean - @ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq", name = "type", havingValue = "kafka", matchIfMissing = true) public CircuitBreaker kafkaSinkCircuitBreaker() { return CircuitBreaker.of("kafka-sink", CircuitBreakerConfig.custom() .slidingWindowSize(100) @@ -61,8 +52,7 @@ public class SinkMqAutoConfiguration { @Bean(destroyMethod = "close") @ConditionalOnMissingBean - @ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq", name = "type", havingValue = "kafka", matchIfMissing = true) - public KafkaProducer kafkaProducer(SinkMqProperties props) { + public KafkaProducer kafkaProducer(KafkaSinkProperties props) { Properties p = new Properties(); p.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, props.getBootstrapServers()); p.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); @@ -77,11 +67,10 @@ public class SinkMqAutoConfiguration { @Bean(destroyMethod = "close") @ConditionalOnMissingBean - @ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq", name = "type", havingValue = "kafka", matchIfMissing = true) public KafkaEventSink kafkaEventSink(KafkaProducer producer, EnvelopeMapper mapper, TopicRouter router, - SinkMqProperties props, + KafkaSinkProperties props, CircuitBreaker breaker) { // KafkaEventSink 是 EventBus 的生产端出口;它不会启动任何 Kafka 消费线程。 return new KafkaEventSink(producer, mapper, router, props.getTopics().getDlq(), breaker); @@ -89,9 +78,8 @@ public class SinkMqAutoConfiguration { @Bean @ConditionalOnMissingBean - @ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq", name = "type", havingValue = "kafka", matchIfMissing = true) public KafkaEnvelopeDeadLetterSink kafkaEnvelopeDeadLetterSink(KafkaProducer producer, - SinkMqProperties props) { + KafkaSinkProperties props) { return new KafkaEnvelopeDeadLetterSink(producer, props.getTopics().getDlq()); } diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqConsumerAutoConfiguration.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfiguration.java similarity index 78% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqConsumerAutoConfiguration.java rename to modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfiguration.java index 025fc10c..602dc8ed 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqConsumerAutoConfiguration.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfiguration.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor; import org.springframework.beans.factory.ListableBeanFactory; @@ -13,15 +13,15 @@ import java.time.Duration; import java.util.List; import java.util.Map; -@AutoConfiguration(after = SinkMqAutoConfiguration.class) +@AutoConfiguration(after = KafkaSinkAutoConfiguration.class) @AutoConfigureAfter(name = { "com.lingniu.ingest.eventhistory.config.EventHistoryAutoConfiguration", "com.lingniu.ingest.vehiclestate.config.VehicleStateAutoConfiguration", "com.lingniu.ingest.vehiclestat.config.VehicleStatAutoConfiguration" }) -@EnableConfigurationProperties(SinkMqProperties.class) -@ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq", name = "enabled", havingValue = "true", matchIfMissing = true) -public class SinkMqConsumerAutoConfiguration { +@EnableConfigurationProperties(KafkaSinkProperties.class) +@ConditionalOnProperty(prefix = "lingniu.ingest.sink.kafka", name = "enabled", havingValue = "true", matchIfMissing = true) +public class KafkaSinkConsumerAutoConfiguration { @Bean @ConditionalOnMissingBean @@ -31,10 +31,10 @@ public class SinkMqConsumerAutoConfiguration { @Bean @ConditionalOnMissingBean - @ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq.consumer", name = "enabled", havingValue = "true") + @ConditionalOnProperty(prefix = "lingniu.ingest.sink.kafka.consumer", name = "enabled", havingValue = "true") public KafkaEnvelopeConsumerRunner kafkaEnvelopeConsumerRunner(ListableBeanFactory beanFactory, KafkaEnvelopeConsumerFactory consumerFactory, - SinkMqProperties props) { + KafkaSinkProperties props) { return new KafkaEnvelopeConsumerRunner( () -> createWorkers(beanFactory, consumerFactory, props), Duration.ofMillis(props.getConsumer().getPollTimeoutMillis()), @@ -44,7 +44,7 @@ public class SinkMqConsumerAutoConfiguration { private List createWorkers(ListableBeanFactory beanFactory, KafkaEnvelopeConsumerFactory consumerFactory, - SinkMqProperties props) { + KafkaSinkProperties props) { Map processors = beanFactory.getBeansOfType(EnvelopeConsumerProcessor.class); return consumerFactory.createWorkers(processors, props); } diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqProperties.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkProperties.java similarity index 90% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqProperties.java rename to modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkProperties.java index d24cf2f7..7d895446 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqProperties.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkProperties.java @@ -1,26 +1,23 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import org.springframework.boot.context.properties.ConfigurationProperties; import java.util.ArrayList; import java.util.LinkedHashMap; import java.util.List; -import java.util.Locale; import java.util.Map; -@ConfigurationProperties(prefix = "lingniu.ingest.sink.mq") -public class SinkMqProperties { +@ConfigurationProperties(prefix = "lingniu.ingest.sink.kafka") +public class KafkaSinkProperties { /** - * MQ Sink 总开关。默认 {@code true}。 - * 设为 {@code false} 时 {@link SinkMqAutoConfiguration} 完全不装配任何 Bean(Kafka Producer、 + * Kafka Sink 总开关。默认 {@code true}。 + * 设为 {@code false} 时 {@link KafkaSinkAutoConfiguration} 完全不装配任何 Bean(Kafka Producer、 * EnvelopeMapper、TopicRouter、KafkaEventSink 都不会创建),ingest-core 的 DisruptorEventBus * 仍然正常运行但没有外部 sink(事件落地到 sink-archive 或 Noop 吞掉)。 */ private boolean enabled = true; - /** MQ 后端类型。生产链路只允许 kafka。 */ - private String type = "kafka"; /** Kafka bootstrap servers;生产环境应通过环境变量覆盖,不建议使用默认开发地址。 */ private String bootstrapServers = "114.55.58.251:9092"; private String compressionType = "zstd"; @@ -35,14 +32,6 @@ public class SinkMqProperties { public boolean isEnabled() { return enabled; } public void setEnabled(boolean enabled) { this.enabled = enabled; } - public String getType() { return type; } - public void setType(String type) { - String value = type == null || type.isBlank() ? "kafka" : type.trim().toLowerCase(Locale.ROOT); - if (!"kafka".equals(value)) { - throw new IllegalStateException("sink.mq.type supports only kafka; configured value: " + type); - } - this.type = value; - } public String getBootstrapServers() { return bootstrapServers; } public void setBootstrapServers(String bootstrapServers) { this.bootstrapServers = bootstrapServers; } public String getCompressionType() { return compressionType; } diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/TopicRouter.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/TopicRouter.java similarity index 84% rename from modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/TopicRouter.java rename to modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/TopicRouter.java index b4a7cbc0..f8c984c0 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/TopicRouter.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/kafka/TopicRouter.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import com.lingniu.ingest.api.event.VehicleEvent; @@ -7,9 +7,9 @@ import com.lingniu.ingest.api.event.VehicleEvent; */ public final class TopicRouter { - private final SinkMqProperties.Topics topics; + private final KafkaSinkProperties.Topics topics; - public TopicRouter(SinkMqProperties.Topics topics) { + public TopicRouter(KafkaSinkProperties.Topics topics) { this.topics = topics; } diff --git a/modules/sinks/sink-mq/src/main/proto/vehicle_envelope.proto b/modules/sinks/sink-mq/src/main/proto/vehicle_envelope.proto index 910ed28b..57acbda1 100644 --- a/modules/sinks/sink-mq/src/main/proto/vehicle_envelope.proto +++ b/modules/sinks/sink-mq/src/main/proto/vehicle_envelope.proto @@ -1,9 +1,9 @@ syntax = "proto3"; -package com.lingniu.ingest.sink.mq.proto; +package com.lingniu.ingest.sink.kafka.proto; option java_multiple_files = true; -option java_package = "com.lingniu.ingest.sink.mq.proto"; +option java_package = "com.lingniu.ingest.sink.kafka.proto"; option java_outer_classname = "VehicleEnvelopeProto"; // 统一消息外壳:所有 Topic 共用此 Envelope,payload 通过 oneof 区分具体事件类型。 diff --git a/modules/sinks/sink-mq/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports b/modules/sinks/sink-mq/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports index 6e83b36d..8c283ae7 100644 --- a/modules/sinks/sink-mq/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports +++ b/modules/sinks/sink-mq/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports @@ -1,2 +1,2 @@ -com.lingniu.ingest.sink.mq.SinkMqAutoConfiguration -com.lingniu.ingest.sink.mq.SinkMqConsumerAutoConfiguration +com.lingniu.ingest.sink.kafka.KafkaSinkAutoConfiguration +com.lingniu.ingest.sink.kafka.KafkaSinkConsumerAutoConfiguration diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/EnvelopeMapperTelemetrySnapshotTest.java b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/EnvelopeMapperTelemetrySnapshotTest.java similarity index 97% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/EnvelopeMapperTelemetrySnapshotTest.java rename to modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/EnvelopeMapperTelemetrySnapshotTest.java index 6e5d806e..5f700f70 100644 --- a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/EnvelopeMapperTelemetrySnapshotTest.java +++ b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/EnvelopeMapperTelemetrySnapshotTest.java @@ -1,10 +1,10 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import com.lingniu.ingest.api.ProtocolId; import com.lingniu.ingest.api.event.RawArchiveKeys; import com.lingniu.ingest.api.event.RealtimePayload; import com.lingniu.ingest.api.event.VehicleEvent; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import org.junit.jupiter.api.Test; import java.time.Instant; diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerFactoryTest.java b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerFactoryTest.java similarity index 95% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerFactoryTest.java rename to modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerFactoryTest.java index 7a4a286f..df0b9ef0 100644 --- a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerFactoryTest.java +++ b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerFactoryTest.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor; import com.lingniu.ingest.api.consumer.EnvelopeIngestResult; @@ -23,7 +23,7 @@ class KafkaEnvelopeConsumerFactoryTest { created.add(props); return new MockConsumer<>(OffsetResetStrategy.EARLIEST); }); - SinkMqProperties props = new SinkMqProperties(); + KafkaSinkProperties props = new KafkaSinkProperties(); props.setBootstrapServers("kafka-1:9092"); EnvelopeConsumerProcessor stateProcessor = processor(); EnvelopeConsumerProcessor statProcessor = processor(); @@ -51,7 +51,7 @@ class KafkaEnvelopeConsumerFactoryTest { created.add(props); return new MockConsumer<>(OffsetResetStrategy.EARLIEST); }); - SinkMqProperties props = new SinkMqProperties(); + KafkaSinkProperties props = new KafkaSinkProperties(); props.getConsumer().setConcurrency(3); EnvelopeConsumerProcessor stateProcessor = processor(); diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerRunnerTest.java b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerRunnerTest.java similarity index 97% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerRunnerTest.java rename to modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerRunnerTest.java index f8cd2c8a..0a4d7f15 100644 --- a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerRunnerTest.java +++ b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerRunnerTest.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import org.junit.jupiter.api.Test; diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerWorkerTest.java b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerWorkerTest.java similarity index 99% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerWorkerTest.java rename to modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerWorkerTest.java index 589686b2..9d1a1b02 100644 --- a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerWorkerTest.java +++ b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeConsumerWorkerTest.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import com.lingniu.ingest.api.consumer.EnvelopeConsumerRecord; import com.lingniu.ingest.api.consumer.EnvelopeBatchIngestor; diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeDeadLetterSinkTest.java b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeDeadLetterSinkTest.java similarity index 98% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeDeadLetterSinkTest.java rename to modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeDeadLetterSinkTest.java index 4c9232ca..94f51d96 100644 --- a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeDeadLetterSinkTest.java +++ b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEnvelopeDeadLetterSinkTest.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import com.lingniu.ingest.api.consumer.EnvelopeDeadLetterRecord; import com.lingniu.ingest.api.consumer.EnvelopeIngestResult; diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEventSinkTest.java b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEventSinkTest.java similarity index 94% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEventSinkTest.java rename to modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEventSinkTest.java index 4fbdd843..2aa77fc9 100644 --- a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEventSinkTest.java +++ b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaEventSinkTest.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import com.lingniu.ingest.api.ProtocolId; import com.lingniu.ingest.api.event.LocationPayload; @@ -26,7 +26,7 @@ class KafkaEventSinkTest { KafkaEventSink sink = new KafkaEventSink( null, new EnvelopeMapper("node-1"), - new TopicRouter(new SinkMqProperties.Topics()), + new TopicRouter(new KafkaSinkProperties.Topics()), "vehicle.dlq.gb32960.v1", CircuitBreaker.ofDefaults("kafka-test")); @@ -46,7 +46,7 @@ class KafkaEventSinkTest { KafkaEventSink sink = new KafkaEventSink( producer, new EnvelopeMapper("node-1"), - new TopicRouter(new SinkMqProperties.Topics()), + new TopicRouter(new KafkaSinkProperties.Topics()), "vehicle.dlq.gb32960.v1", CircuitBreaker.ofDefaults("kafka-test")); diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/SinkMqConsumerAutoConfigurationTest.java b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfigurationTest.java similarity index 70% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/SinkMqConsumerAutoConfigurationTest.java rename to modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfigurationTest.java index 72126b59..25b570ea 100644 --- a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/SinkMqConsumerAutoConfigurationTest.java +++ b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfigurationTest.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor; import com.lingniu.ingest.api.consumer.EnvelopeIngestResult; @@ -15,26 +15,25 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.Mockito.mock; @ExtendWith(OutputCaptureExtension.class) -class SinkMqConsumerAutoConfigurationTest { +class KafkaSinkConsumerAutoConfigurationTest { private final ApplicationContextRunner contextRunner = new ApplicationContextRunner() - .withUserConfiguration(SinkMqAutoConfiguration.class, SinkMqConsumerAutoConfiguration.class) + .withUserConfiguration(KafkaSinkAutoConfiguration.class, KafkaSinkConsumerAutoConfiguration.class) .withAllowBeanDefinitionOverriding(true) .withBean("vehicleStateEnvelopeConsumerProcessor", EnvelopeConsumerProcessor.class, () -> new EnvelopeConsumerProcessor( "vehicle-state", bytes -> EnvelopeIngestResult.processed("evt-1", "VIN001"), record -> {})) - .withBean("kafkaProducer", KafkaProducer.class, SinkMqConsumerAutoConfigurationTest::kafkaProducer) + .withBean("kafkaProducer", KafkaProducer.class, KafkaSinkConsumerAutoConfigurationTest::kafkaProducer) .withBean(KafkaEnvelopeConsumerFactory.class, () -> new KafkaEnvelopeConsumerFactory(props -> new MockConsumer<>(OffsetResetStrategy.EARLIEST))) .withPropertyValues( - "lingniu.ingest.sink.mq.enabled=true", - "lingniu.ingest.sink.mq.type=kafka", - "lingniu.ingest.sink.mq.consumer.enabled=true", - "lingniu.ingest.sink.mq.consumer.auto-startup=false", - "lingniu.ingest.sink.mq.consumer.bindings.vehicleStateEnvelopeConsumerProcessor.group-id=vehicle-state", - "lingniu.ingest.sink.mq.consumer.bindings.vehicleStateEnvelopeConsumerProcessor.topics[0]=vehicle.realtime"); + "lingniu.ingest.sink.kafka.enabled=true", + "lingniu.ingest.sink.kafka.consumer.enabled=true", + "lingniu.ingest.sink.kafka.consumer.auto-startup=false", + "lingniu.ingest.sink.kafka.consumer.bindings.vehicleStateEnvelopeConsumerProcessor.group-id=vehicle-state", + "lingniu.ingest.sink.kafka.consumer.bindings.vehicleStateEnvelopeConsumerProcessor.topics[0]=vehicle.realtime"); @Test void createsKafkaEnvelopeConsumerRunnerWhenConsumerBindingIsConfigured(CapturedOutput output) { diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkPropertiesTest.java b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkPropertiesTest.java new file mode 100644 index 00000000..dbe4ab4c --- /dev/null +++ b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkPropertiesTest.java @@ -0,0 +1,26 @@ +package com.lingniu.ingest.sink.kafka; + +import org.junit.jupiter.api.Test; + +import java.util.Arrays; + +import static org.assertj.core.api.Assertions.assertThat; + +class KafkaSinkPropertiesTest { + + @Test + void defaultKafkaBrokerUsesProductionAddressButCanStillBeOverriddenByConfigBinding() { + KafkaSinkProperties props = new KafkaSinkProperties(); + + assertThat(props.getBootstrapServers()).isEqualTo("114.55.58.251:9092"); + + props.setBootstrapServers("kafka.internal:9092"); + assertThat(props.getBootstrapServers()).isEqualTo("kafka.internal:9092"); + } + + @Test + void exposesNoBackendTypeBecauseKafkaIsTheOnlySink() { + assertThat(Arrays.stream(KafkaSinkProperties.class.getMethods()).map(method -> method.getName())) + .doesNotContain("getType", "setType"); + } +} diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/TopicRouterTest.java b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/TopicRouterTest.java similarity index 94% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/TopicRouterTest.java rename to modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/TopicRouterTest.java index 7ba6c418..7642426d 100644 --- a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/TopicRouterTest.java +++ b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/TopicRouterTest.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; import com.lingniu.ingest.api.ProtocolId; import com.lingniu.ingest.api.event.AlarmPayload; @@ -18,14 +18,14 @@ class TopicRouterTest { @Test void rawArchiveRoutesToVersionedGb32960RawTopicByDefault() { - TopicRouter router = new TopicRouter(new SinkMqProperties.Topics()); + TopicRouter router = new TopicRouter(new KafkaSinkProperties.Topics()); assertThat(router.route(rawArchive())).isEqualTo("vehicle.raw.gb32960.v1"); } @Test void normalizedGb32960EventsRouteToVersionedEventTopicByDefault() { - TopicRouter router = new TopicRouter(new SinkMqProperties.Topics()); + TopicRouter router = new TopicRouter(new KafkaSinkProperties.Topics()); assertThat(router.route(realtime())).isEqualTo("vehicle.event.gb32960.v1"); assertThat(router.route(location())).isEqualTo("vehicle.event.gb32960.v1"); @@ -37,7 +37,7 @@ class TopicRouterTest { @Test void dlqDefaultsToVersionedGb32960DlqTopic() { - SinkMqProperties.Topics topics = new SinkMqProperties.Topics(); + KafkaSinkProperties.Topics topics = new KafkaSinkProperties.Topics(); assertThat(topics.getDlq()).isEqualTo("vehicle.dlq.gb32960.v1"); } diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/VehicleEnvelopeProtoCompatibilityTest.java b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/VehicleEnvelopeProtoCompatibilityTest.java similarity index 90% rename from modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/VehicleEnvelopeProtoCompatibilityTest.java rename to modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/VehicleEnvelopeProtoCompatibilityTest.java index c68d5a76..405031e2 100644 --- a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/VehicleEnvelopeProtoCompatibilityTest.java +++ b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/kafka/VehicleEnvelopeProtoCompatibilityTest.java @@ -1,9 +1,9 @@ -package com.lingniu.ingest.sink.mq; +package com.lingniu.ingest.sink.kafka; -import com.lingniu.ingest.sink.mq.proto.DecodedFactPayload; -import com.lingniu.ingest.sink.mq.proto.ParseStatusProto; -import com.lingniu.ingest.sink.mq.proto.RawFrameFactPayload; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.DecodedFactPayload; +import com.lingniu.ingest.sink.kafka.proto.ParseStatusProto; +import com.lingniu.ingest.sink.kafka.proto.RawFrameFactPayload; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import org.junit.jupiter.api.Test; import static org.assertj.core.api.Assertions.assertThat; diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/SinkMqPropertiesTest.java b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/SinkMqPropertiesTest.java deleted file mode 100644 index 4f3c54dd..00000000 --- a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/SinkMqPropertiesTest.java +++ /dev/null @@ -1,29 +0,0 @@ -package com.lingniu.ingest.sink.mq; - -import org.junit.jupiter.api.Test; - -import static org.assertj.core.api.Assertions.assertThat; -import static org.assertj.core.api.Assertions.assertThatThrownBy; - -class SinkMqPropertiesTest { - - @Test - void defaultKafkaBrokerUsesProductionAddressButCanStillBeOverriddenByConfigBinding() { - SinkMqProperties props = new SinkMqProperties(); - - assertThat(props.getBootstrapServers()).isEqualTo("114.55.58.251:9092"); - - props.setBootstrapServers("kafka.internal:9092"); - assertThat(props.getBootstrapServers()).isEqualTo("kafka.internal:9092"); - } - - @Test - void rejectsUnsupportedMqTypeInsteadOfSilentlyDisablingKafkaBeans() { - SinkMqProperties props = new SinkMqProperties(); - - assertThatThrownBy(() -> props.setType("rocketmq")) - .isInstanceOf(IllegalStateException.class) - .hasMessageContaining("sink.mq.type supports only kafka") - .hasMessageContaining("rocketmq"); - } -} diff --git a/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineEnvelopeRows.java b/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineEnvelopeRows.java index 3119a079..9d654148 100644 --- a/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineEnvelopeRows.java +++ b/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineEnvelopeRows.java @@ -2,10 +2,10 @@ package com.lingniu.ingest.tdenginehistory; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; -import com.lingniu.ingest.sink.mq.proto.LocationPayload; -import com.lingniu.ingest.sink.mq.proto.RawFrameFactPayload; -import com.lingniu.ingest.sink.mq.proto.TelemetryField; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.LocationPayload; +import com.lingniu.ingest.sink.kafka.proto.RawFrameFactPayload; +import com.lingniu.ingest.sink.kafka.proto.TelemetryField; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import java.time.Instant; import java.util.ArrayList; diff --git a/modules/sinks/tdengine-history-store/src/test/java/com/lingniu/ingest/tdenginehistory/TdengineEnvelopeRowsTest.java b/modules/sinks/tdengine-history-store/src/test/java/com/lingniu/ingest/tdenginehistory/TdengineEnvelopeRowsTest.java index e70e40b9..023b8dfd 100644 --- a/modules/sinks/tdengine-history-store/src/test/java/com/lingniu/ingest/tdenginehistory/TdengineEnvelopeRowsTest.java +++ b/modules/sinks/tdengine-history-store/src/test/java/com/lingniu/ingest/tdenginehistory/TdengineEnvelopeRowsTest.java @@ -1,12 +1,12 @@ package com.lingniu.ingest.tdenginehistory; -import com.lingniu.ingest.sink.mq.proto.LocationPayload; -import com.lingniu.ingest.sink.mq.proto.ParseStatusProto; -import com.lingniu.ingest.sink.mq.proto.RawArchiveRef; -import com.lingniu.ingest.sink.mq.proto.RawFrameFactPayload; -import com.lingniu.ingest.sink.mq.proto.TelemetryField; -import com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot; -import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.kafka.proto.LocationPayload; +import com.lingniu.ingest.sink.kafka.proto.ParseStatusProto; +import com.lingniu.ingest.sink.kafka.proto.RawArchiveRef; +import com.lingniu.ingest.sink.kafka.proto.RawFrameFactPayload; +import com.lingniu.ingest.sink.kafka.proto.TelemetryField; +import com.lingniu.ingest.sink.kafka.proto.TelemetrySnapshot; +import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope; import org.junit.jupiter.api.Test; import java.time.Instant;