From 3161e98eff217276621beba59e0c321358155ace Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 15:01:38 +0800 Subject: [PATCH] chore: scope kafka consumer env to consumers --- deploy/portainer/docker-compose.yml | 9 ++++--- .../src/main/resources/application.yml | 2 +- .../PortainerComposeResourceLimitsTest.java | 24 +++++++++++++++++++ .../VehicleHistoryAppDefaultsTest.java | 2 +- 4 files changed, 30 insertions(+), 7 deletions(-) diff --git a/deploy/portainer/docker-compose.yml b/deploy/portainer/docker-compose.yml index c8bca16b..6133275d 100644 --- a/deploy/portainer/docker-compose.yml +++ b/deploy/portainer/docker-compose.yml @@ -17,8 +17,6 @@ x-common-env: &common-env NACOS_USERNAME: ${NACOS_USERNAME:-} NACOS_PASSWORD: ${NACOS_PASSWORD:-} KAFKA_BROKERS: ${KAFKA_BROKERS:-172.17.111.56:9092} - KAFKA_CONSUMER_ENABLED: ${KAFKA_CONSUMER_ENABLED:-true} - KAFKA_CONSUMER_MAX_POLL_INTERVAL_MILLIS: ${KAFKA_CONSUMER_MAX_POLL_INTERVAL_MILLIS:-900000} KAFKA_TOPIC_GB32960_EVENT: ${KAFKA_TOPIC_GB32960_EVENT:-vehicle.event.gb32960.v1} KAFKA_TOPIC_GB32960_RAW: ${KAFKA_TOPIC_GB32960_RAW:-vehicle.raw.gb32960.v1} KAFKA_TOPIC_GB32960_DLQ: ${KAFKA_TOPIC_GB32960_DLQ:-vehicle.dlq.gb32960.v1} @@ -63,7 +61,6 @@ services: - *redis-session-env HTTP_PORT: 20100 GB32960_PORT: 32960 - KAFKA_CONSUMER_ENABLED: "false" KAFKA_NODE_ID: ${KAFKA_NODE_ID:-gb32960-ingest-portainer} GB32960_AUTH_ENABLED: ${GB32960_AUTH_ENABLED:-false} GB32960_PLATFORM_USER_HYUNDAI: ${GB32960_PLATFORM_USER_HYUNDAI:-Hyundai} @@ -97,7 +94,6 @@ services: - *redis-session-env HTTP_PORT: 20400 JT808_PORT: 808 - KAFKA_CONSUMER_ENABLED: "false" KAFKA_NODE_ID: ${KAFKA_NODE_ID_JT808:-jt808-ingest-portainer} SINK_ARCHIVE_PATH: ${JT808_ARCHIVE_PATH:-/data/archive} ports: @@ -123,7 +119,6 @@ services: HTTP_PORT: 20500 YUTONG_MQTT_ENABLED: ${YUTONG_MQTT_ENABLED:-false} YUTONG_MQTT_AUTO_STARTUP: ${YUTONG_MQTT_AUTO_STARTUP:-true} - KAFKA_CONSUMER_ENABLED: "false" KAFKA_NODE_ID: ${KAFKA_NODE_ID_MQTT:-yutong-mqtt-portainer} YUTONG_MQTT_ENDPOINT_NAME: ${YUTONG_MQTT_ENDPOINT_NAME:-yutong} YUTONG_MQTT_URI: ${YUTONG_MQTT_URI:-} @@ -162,8 +157,10 @@ services: environment: <<: *common-env HTTP_PORT: 20200 + KAFKA_CONSUMER_ENABLED: ${KAFKA_CONSUMER_ENABLED_HISTORY:-true} KAFKA_CONSUMER_CLIENT_ID_PREFIX: ${KAFKA_CONSUMER_CLIENT_ID_PREFIX_HISTORY:-vehicle-history} KAFKA_CONSUMER_MAX_POLL_RECORDS: ${KAFKA_CONSUMER_MAX_POLL_RECORDS_HISTORY:-2000} + KAFKA_CONSUMER_MAX_POLL_INTERVAL_MILLIS: ${KAFKA_CONSUMER_MAX_POLL_INTERVAL_MILLIS_HISTORY:-1800000} KAFKA_GROUP_HISTORY: ${KAFKA_GROUP_HISTORY:-vehicle-history} KAFKA_GROUP_HISTORY_GB32960_EVENT: ${KAFKA_GROUP_HISTORY_GB32960_EVENT:-${KAFKA_GROUP_HISTORY:-vehicle-history}-gb32960-event} KAFKA_GROUP_HISTORY_JT808_EVENT: ${KAFKA_GROUP_HISTORY_JT808_EVENT:-${KAFKA_GROUP_HISTORY:-vehicle-history}-jt808-event} @@ -197,7 +194,9 @@ services: environment: <<: *common-env HTTP_PORT: 20300 + KAFKA_CONSUMER_ENABLED: ${KAFKA_CONSUMER_ENABLED_ANALYTICS:-true} KAFKA_CONSUMER_CLIENT_ID_PREFIX: ${KAFKA_CONSUMER_CLIENT_ID_PREFIX_ANALYTICS:-vehicle-analytics} + KAFKA_CONSUMER_MAX_POLL_INTERVAL_MILLIS: ${KAFKA_CONSUMER_MAX_POLL_INTERVAL_MILLIS_ANALYTICS:-900000} KAFKA_GROUP_STAT: ${KAFKA_GROUP_STAT:-vehicle-stat} VEHICLE_STAT_ENABLED: ${VEHICLE_STAT_ENABLED:-true} VEHICLE_STAT_ZONE_ID: ${VEHICLE_STAT_ZONE_ID:-Asia/Shanghai} 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 648d4007..2c8b25b4 100644 --- a/modules/apps/vehicle-history-app/src/main/resources/application.yml +++ b/modules/apps/vehicle-history-app/src/main/resources/application.yml @@ -46,7 +46,7 @@ lingniu: client-id-prefix: ${KAFKA_CONSUMER_CLIENT_ID_PREFIX:vehicle-history} auto-offset-reset: ${KAFKA_CONSUMER_AUTO_OFFSET_RESET:earliest} max-poll-records: ${KAFKA_CONSUMER_MAX_POLL_RECORDS:2000} - max-poll-interval-millis: ${KAFKA_CONSUMER_MAX_POLL_INTERVAL_MS:1800000} + max-poll-interval-millis: ${KAFKA_CONSUMER_MAX_POLL_INTERVAL_MILLIS:1800000} concurrency: ${KAFKA_CONSUMER_CONCURRENCY:3} bindings: eventHistoryGb32960EnvelopeConsumerProcessor: 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 b71c0707..9b755964 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 @@ -309,6 +309,30 @@ class PortainerComposeResourceLimitsTest { .doesNotContain("\n KAFKA_GROUP_HISTORY_MQTT_YUTONG:"); } + @Test + void kafkaConsumerEnvironmentIsScopedToConsumerAppsOnly() throws IOException { + String compose = Files.readString(repositoryRoot().resolve("deploy/portainer/docker-compose.yml")); + String commonEnv = commonEnvBlock(compose); + String ingestServices = serviceBlock(compose, "gb32960-ingest-app") + + serviceBlock(compose, "jt808-ingest-app") + + serviceBlock(compose, "yutong-mqtt-app"); + String historyService = serviceBlock(compose, "vehicle-history-app"); + String analyticsService = serviceBlock(compose, "vehicle-analytics-app"); + + assertThat(commonEnv) + .doesNotContain("KAFKA_CONSUMER_ENABLED") + .doesNotContain("KAFKA_CONSUMER_MAX_POLL_INTERVAL"); + assertThat(ingestServices) + .doesNotContain("KAFKA_CONSUMER_"); + assertThat(historyService) + .contains("KAFKA_CONSUMER_ENABLED: ${KAFKA_CONSUMER_ENABLED_HISTORY:-true}") + .contains("KAFKA_CONSUMER_MAX_POLL_INTERVAL_MILLIS: ${KAFKA_CONSUMER_MAX_POLL_INTERVAL_MILLIS_HISTORY:-1800000}") + .doesNotContain("KAFKA_CONSUMER_MAX_POLL_INTERVAL_MS"); + assertThat(analyticsService) + .contains("KAFKA_CONSUMER_ENABLED: ${KAFKA_CONSUMER_ENABLED_ANALYTICS:-true}") + .contains("KAFKA_CONSUMER_MAX_POLL_INTERVAL_MILLIS: ${KAFKA_CONSUMER_MAX_POLL_INTERVAL_MILLIS_ANALYTICS:-900000}"); + } + @Test void vehicleHistoryUsesTdengineHistoryEnvironmentNames() throws IOException { String compose = Files.readString(repositoryRoot().resolve("deploy/portainer/docker-compose.yml")); 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 71589975..89f87754 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 @@ -46,7 +46,7 @@ class VehicleHistoryAppDefaultsTest { .containsEntry("lingniu.ingest.sink.kafka.consumer.concurrency", "${KAFKA_CONSUMER_CONCURRENCY:3}") .containsEntry( "lingniu.ingest.sink.kafka.consumer.max-poll-interval-millis", - "${KAFKA_CONSUMER_MAX_POLL_INTERVAL_MS:1800000}") + "${KAFKA_CONSUMER_MAX_POLL_INTERVAL_MILLIS:1800000}") .containsEntry( "lingniu.ingest.sink.kafka.consumer.bindings.eventHistoryGb32960EnvelopeConsumerProcessor.enabled", true)