chore: scope kafka consumer env to consumers

This commit is contained in:
lingniu
2026-07-01 15:01:38 +08:00
parent 5910d23dd9
commit 3161e98eff
4 changed files with 30 additions and 7 deletions

View File

@@ -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}

View File

@@ -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:

View File

@@ -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"));

View File

@@ -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)