From 6ac77b2f7d26ac19cfbaac87097529cf13e4b047 Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 16:30:51 +0800 Subject: [PATCH] chore: separate kafka consumer dlq topic --- deploy/portainer/docker-compose.yml | 1 + .../gb32960-service-split-runbook.md | 4 +++- .../src/main/resources/application.yml | 4 +--- .../VehicleAnalyticsAppDefaultsTest.java | 3 +++ .../src/main/resources/application.yml | 1 + .../PortainerComposeResourceLimitsTest.java | 2 ++ .../VehicleHistoryAppDefaultsTest.java | 2 ++ .../kafka/KafkaSinkAutoConfiguration.java | 2 +- .../sink/kafka/KafkaSinkProperties.java | 12 ++++++++++ ...afkaSinkConsumerAutoConfigurationTest.java | 23 +++++++++++++++++++ 10 files changed, 49 insertions(+), 5 deletions(-) diff --git a/deploy/portainer/docker-compose.yml b/deploy/portainer/docker-compose.yml index 13c1fd6e..6d497567 100644 --- a/deploy/portainer/docker-compose.yml +++ b/deploy/portainer/docker-compose.yml @@ -160,6 +160,7 @@ services: KAFKA_TOPIC_JT808_RAW: ${KAFKA_TOPIC_JT808_RAW:-vehicle.raw.jt808.v1} KAFKA_TOPIC_YUTONG_MQTT_EVENT: ${KAFKA_TOPIC_YUTONG_MQTT_EVENT:-vehicle.event.mqtt-yutong.v1} KAFKA_TOPIC_YUTONG_MQTT_RAW: ${KAFKA_TOPIC_YUTONG_MQTT_RAW:-vehicle.raw.mqtt-yutong.v1} + KAFKA_TOPIC_HISTORY_DLQ: ${KAFKA_TOPIC_HISTORY_DLQ:-vehicle.dlq.history.v1} KAFKA_GROUP_HISTORY: ${KAFKA_GROUP_HISTORY:-vehicle-history} TDENGINE_HISTORY_ENABLED: ${TDENGINE_HISTORY_ENABLED:-true} TDENGINE_JDBC_URL: ${TDENGINE_JDBC_URL:-jdbc:TAOS-WS://172.17.111.57:6041/vehicle_ts} diff --git a/docs/operations/gb32960-service-split-runbook.md b/docs/operations/gb32960-service-split-runbook.md index ead061ab..1c16384d 100644 --- a/docs/operations/gb32960-service-split-runbook.md +++ b/docs/operations/gb32960-service-split-runbook.md @@ -25,6 +25,7 @@ This runbook covers the default production runtimes: GB32960 ingest, JT808 inges | `vehicle.raw.mqtt-yutong.v1` | `yutong-mqtt-app` | `vehicle-history-app` | Yutong MQTT raw payload records for TDengine `raw_frames` and troubleshooting APIs. | | `vehicle.event.mqtt-yutong.v1` | `yutong-mqtt-app` | `vehicle-history-app` | Yutong MQTT parsed telemetry events keyed by mapped VIN or device identity. | | `vehicle.dlq.mqtt-yutong.v1` | Kafka producer/consumer error paths | Operators/replay tooling | Dead-letter records for failed Yutong MQTT production or consumer processing. | +| `vehicle.dlq.history.v1` | History app consumer error paths | Operators/replay tooling | Dead-letter records for failed history storage processing. | Default local consumer groups: @@ -84,12 +85,13 @@ kafka-topics --bootstrap-server 127.0.0.1:9092 --create --if-not-exists --topic kafka-topics --bootstrap-server 127.0.0.1:9092 --create --if-not-exists --topic vehicle.raw.mqtt-yutong.v1 --partitions 12 --replication-factor 1 kafka-topics --bootstrap-server 127.0.0.1:9092 --create --if-not-exists --topic vehicle.event.mqtt-yutong.v1 --partitions 12 --replication-factor 1 kafka-topics --bootstrap-server 127.0.0.1:9092 --create --if-not-exists --topic vehicle.dlq.mqtt-yutong.v1 --partitions 3 --replication-factor 1 +kafka-topics --bootstrap-server 127.0.0.1:9092 --create --if-not-exists --topic vehicle.dlq.history.v1 --partitions 3 --replication-factor 1 ``` Optional sanity check: ```bash -kafka-topics --bootstrap-server 127.0.0.1:9092 --list | grep -E 'vehicle\.(raw|event|dlq)\.(gb32960|jt808|mqtt-yutong)\.v1' +kafka-topics --bootstrap-server 127.0.0.1:9092 --list | grep -E 'vehicle\.(raw|event|dlq)\.(gb32960|jt808|mqtt-yutong|history)\.v1' ``` ## Start Split Services Locally 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 25b455ba..b5fac976 100644 --- a/modules/apps/vehicle-analytics-app/src/main/resources/application.yml +++ b/modules/apps/vehicle-analytics-app/src/main/resources/application.yml @@ -38,10 +38,8 @@ lingniu: kafka: enabled: true bootstrap-servers: ${KAFKA_BROKERS:114.55.58.251:9092} - topics: - realtime: ${KAFKA_TOPIC_JT808_EVENT:vehicle.event.jt808.v1} - dlq: ${KAFKA_TOPIC_JT808_DLQ:vehicle.dlq.jt808.v1} consumer: + dlq-topic: ${KAFKA_TOPIC_JT808_DLQ:vehicle.dlq.jt808.v1} enabled: ${KAFKA_CONSUMER_ENABLED:true} client-id-prefix: ${KAFKA_CONSUMER_CLIENT_ID_PREFIX:vehicle-analytics} auto-offset-reset: ${KAFKA_CONSUMER_AUTO_OFFSET_RESET:earliest} 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 f87e64be..34909d50 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 @@ -28,6 +28,8 @@ class VehicleAnalyticsAppDefaultsTest { "${VEHICLE_STAT_JT808_MILEAGE_ENABLED:true}") .containsEntry("lingniu.ingest.sink.kafka.enabled", true) .containsEntry("lingniu.ingest.sink.kafka.consumer.enabled", "${KAFKA_CONSUMER_ENABLED:true}") + .containsEntry("lingniu.ingest.sink.kafka.consumer.dlq-topic", + "${KAFKA_TOPIC_JT808_DLQ:vehicle.dlq.jt808.v1}") .containsEntry( "lingniu.ingest.sink.kafka.consumer.bindings.vehicleStatEnvelopeConsumerProcessor.enabled", true) @@ -40,6 +42,7 @@ class VehicleAnalyticsAppDefaultsTest { .noneMatch(name -> name.startsWith("lingniu.ingest.gb32960.")) .noneMatch(name -> name.startsWith("lingniu.ingest.event-history.")) .noneMatch(name -> name.startsWith("lingniu.ingest.sink.archive.")) + .noneMatch(name -> name.startsWith("lingniu.ingest.sink.kafka.topics.")) .noneMatch(name -> name.startsWith("lingniu.ingest.event-file-store.")) .noneMatch(name -> name.startsWith("lingniu.ingest.vehicle-state.")) .noneMatch(name -> name.contains("vehicleStateEnvelopeConsumerProcessor")) 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 88e2c1d5..0b08b0de 100644 --- a/modules/apps/vehicle-history-app/src/main/resources/application.yml +++ b/modules/apps/vehicle-history-app/src/main/resources/application.yml @@ -34,6 +34,7 @@ lingniu: enabled: true bootstrap-servers: ${KAFKA_BROKERS:114.55.58.251:9092} consumer: + dlq-topic: ${KAFKA_TOPIC_HISTORY_DLQ:vehicle.dlq.history.v1} enabled: ${KAFKA_CONSUMER_ENABLED:true} client-id-prefix: ${KAFKA_CONSUMER_CLIENT_ID_PREFIX:vehicle-history} auto-offset-reset: ${KAFKA_CONSUMER_AUTO_OFFSET_RESET:earliest} diff --git a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/PortainerComposeResourceLimitsTest.java b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/PortainerComposeResourceLimitsTest.java index e6770e18..88df5662 100644 --- a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/PortainerComposeResourceLimitsTest.java +++ b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/PortainerComposeResourceLimitsTest.java @@ -369,6 +369,7 @@ class PortainerComposeResourceLimitsTest { .contains("KAFKA_TOPIC_JT808_RAW: ${KAFKA_TOPIC_JT808_RAW:-vehicle.raw.jt808.v1}") .contains("KAFKA_TOPIC_YUTONG_MQTT_EVENT: ${KAFKA_TOPIC_YUTONG_MQTT_EVENT:-vehicle.event.mqtt-yutong.v1}") .contains("KAFKA_TOPIC_YUTONG_MQTT_RAW: ${KAFKA_TOPIC_YUTONG_MQTT_RAW:-vehicle.raw.mqtt-yutong.v1}") + .contains("KAFKA_TOPIC_HISTORY_DLQ: ${KAFKA_TOPIC_HISTORY_DLQ:-vehicle.dlq.history.v1}") .doesNotContain("KAFKA_TOPIC_GB32960_DLQ") .doesNotContain("KAFKA_TOPIC_JT808_DLQ") .doesNotContain("KAFKA_TOPIC_YUTONG_MQTT_DLQ"); @@ -493,6 +494,7 @@ class PortainerComposeResourceLimitsTest { .contains("`vehicle.raw.mqtt-yutong.v1`") .contains("`vehicle.event.mqtt-yutong.v1`") .contains("`vehicle.dlq.mqtt-yutong.v1`") + .contains("`vehicle.dlq.history.v1`") .contains(activeBuildCommand) .contains("modules/apps/yutong-mqtt-app/target/yutong-mqtt-app.jar") .doesNotContain("20482") 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 1a9432d2..313bd745 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 @@ -40,6 +40,8 @@ class VehicleHistoryAppDefaultsTest { .containsEntry("lingniu.ingest.event-history.enabled", true) .containsEntry("lingniu.ingest.sink.kafka.enabled", true) .containsEntry("lingniu.ingest.sink.kafka.consumer.enabled", "${KAFKA_CONSUMER_ENABLED:true}") + .containsEntry("lingniu.ingest.sink.kafka.consumer.dlq-topic", + "${KAFKA_TOPIC_HISTORY_DLQ:vehicle.dlq.history.v1}") .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( diff --git a/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkAutoConfiguration.java b/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkAutoConfiguration.java index af7ff127..82a1dd91 100644 --- a/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkAutoConfiguration.java +++ b/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkAutoConfiguration.java @@ -80,7 +80,7 @@ public class KafkaSinkAutoConfiguration { @ConditionalOnMissingBean public KafkaEnvelopeDeadLetterSink kafkaEnvelopeDeadLetterSink(KafkaProducer producer, KafkaSinkProperties props) { - return new KafkaEnvelopeDeadLetterSink(producer, props.getTopics().getDlq()); + return new KafkaEnvelopeDeadLetterSink(producer, props.deadLetterTopic()); } } diff --git a/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkProperties.java b/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkProperties.java index 8296ef07..69d55521 100644 --- a/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkProperties.java +++ b/modules/sinks/sink-kafka/src/main/java/com/lingniu/ingest/sink/kafka/KafkaSinkProperties.java @@ -51,6 +51,14 @@ public class KafkaSinkProperties { public Consumer getConsumer() { return consumer; } public void setConsumer(Consumer consumer) { this.consumer = consumer; } + public String deadLetterTopic() { + String consumerDlq = consumer == null ? "" : consumer.getDlqTopic(); + if (consumerDlq != null && !consumerDlq.isBlank()) { + return consumerDlq; + } + return topics.getDlq(); + } + public static class Topics { /** 实时遥测事件 topic;GB32960 RAW-only 架构下可逐步弱化该 topic。 */ private String realtime = "vehicle.event.gb32960.v1"; @@ -91,6 +99,8 @@ public class KafkaSinkProperties { private int maxPollRecords = 500; private int maxPollIntervalMillis = 1800000; private int concurrency = 1; + /** 消费处理失败时的死信 topic;未配置时兼容使用 producer topics.dlq。 */ + private String dlqTopic = ""; /** * Processor bean name -> Kafka binding。 * @@ -117,6 +127,8 @@ public class KafkaSinkProperties { public void setMaxPollIntervalMillis(int maxPollIntervalMillis) { this.maxPollIntervalMillis = maxPollIntervalMillis; } public int getConcurrency() { return concurrency; } public void setConcurrency(int concurrency) { this.concurrency = concurrency; } + public String getDlqTopic() { return dlqTopic; } + public void setDlqTopic(String dlqTopic) { this.dlqTopic = dlqTopic; } public Map getBindings() { return bindings; } public void setBindings(Map bindings) { this.bindings = bindings; } } diff --git a/modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfigurationTest.java b/modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfigurationTest.java index 25b570ea..b4d680c5 100644 --- a/modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfigurationTest.java +++ b/modules/sinks/sink-kafka/src/test/java/com/lingniu/ingest/sink/kafka/KafkaSinkConsumerAutoConfigurationTest.java @@ -11,6 +11,8 @@ import org.springframework.boot.test.context.runner.ApplicationContextRunner; import org.springframework.boot.test.system.CapturedOutput; import org.springframework.boot.test.system.OutputCaptureExtension; +import java.lang.reflect.Field; + import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.Mockito.mock; @@ -44,8 +46,29 @@ class KafkaSinkConsumerAutoConfigurationTest { assertThat(output).doesNotContain("114.55.58.251"); } + @Test + void consumerDlqTopicOverridesProducerTopicDefaultsForDeadLetterSink() { + contextRunner + .withPropertyValues("lingniu.ingest.sink.kafka.consumer.dlq-topic=vehicle.dlq.consumer.v1") + .run(context -> { + assertThat(context).hasSingleBean(KafkaEnvelopeDeadLetterSink.class); + assertThat(deadLetterTopic(context.getBean(KafkaEnvelopeDeadLetterSink.class))) + .isEqualTo("vehicle.dlq.consumer.v1"); + }); + } + @SuppressWarnings("unchecked") private static KafkaProducer kafkaProducer() { return mock(KafkaProducer.class); } + + private static String deadLetterTopic(KafkaEnvelopeDeadLetterSink sink) { + try { + Field field = KafkaEnvelopeDeadLetterSink.class.getDeclaredField("topic"); + field.setAccessible(true); + return (String) field.get(sink); + } catch (ReflectiveOperationException e) { + throw new AssertionError(e); + } + } }