refactor: store runtime identity and mileage in mysql metrics

This commit is contained in:
lingniu
2026-07-01 01:56:16 +08:00
parent b1190cd7c7
commit 0fc91f512c
23 changed files with 75 additions and 132 deletions

View File

@@ -25,6 +25,10 @@
<groupId>com.lingniu.ingest</groupId>
<artifactId>observability</artifactId>
</dependency>
<dependency>
<groupId>com.lingniu.ingest</groupId>
<artifactId>vehicle-identity</artifactId>
</dependency>
<dependency>
<groupId>com.lingniu.ingest</groupId>
<artifactId>protocol-gb32960</artifactId>

View File

@@ -78,9 +78,13 @@ lingniu:
store: ${SESSION_STORE:memory}
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:
enabled: ${KAFKA_ENABLED:true}

View File

@@ -1,6 +1,8 @@
package com.lingniu.ingest.gb32960app;
import com.lingniu.ingest.core.config.IngestCoreAutoConfiguration;
import com.lingniu.ingest.identity.MySqlVehicleIdentityService;
import com.lingniu.ingest.identity.config.VehicleIdentityAutoConfiguration;
import com.lingniu.ingest.protocol.gb32960.codec.Gb32960MessageDecoder;
import com.lingniu.ingest.protocol.gb32960.config.Gb32960AutoConfiguration;
import com.lingniu.ingest.protocol.gb32960.inbound.Gb32960NettyServer;
@@ -27,6 +29,7 @@ class Gb32960IngestAppCompositionTest {
.withConfiguration(AutoConfigurations.of(
IngestCoreAutoConfiguration.class,
SessionCoreAutoConfiguration.class,
VehicleIdentityAutoConfiguration.class,
SinkMqAutoConfiguration.class,
SinkArchiveAutoConfiguration.class,
Gb32960AutoConfiguration.class))
@@ -36,6 +39,8 @@ class Gb32960IngestAppCompositionTest {
"lingniu.ingest.gb32960.enabled=true",
"lingniu.ingest.gb32960.server.enabled=true",
"lingniu.ingest.gb32960.port=0",
"lingniu.ingest.identity.store=mysql",
"lingniu.ingest.identity.mysql.initialize-schema=false",
"lingniu.ingest.session.store=memory",
"lingniu.ingest.sink.mq.enabled=true",
"lingniu.ingest.sink.mq.type=kafka",
@@ -53,6 +58,7 @@ class Gb32960IngestAppCompositionTest {
contextRunner.run(context -> {
assertThat(context).hasSingleBean(Gb32960MessageDecoder.class);
assertThat(context).hasSingleBean(Gb32960NettyServer.class);
assertThat(context).hasSingleBean(MySqlVehicleIdentityService.class);
assertThat(context).hasSingleBean(KafkaEventSink.class);
assertThat(context).hasSingleBean(KafkaEnvelopeDeadLetterSink.class);
assertThat(context).hasSingleBean(ArchiveStore.class);

View File

@@ -27,6 +27,9 @@ class Gb32960IngestAppDefaultsTest {
.containsEntry("lingniu.ingest.sink.mq.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.identity.store", "${VEHICLE_IDENTITY_STORE:mysql}")
.containsEntry("lingniu.ingest.identity.mysql.table-name",
"${VEHICLE_IDENTITY_MYSQL_TABLE:vehicle_identity_bindings}")
.containsEntry("lingniu.ingest.event-file-store.enabled", false)
.containsEntry("lingniu.ingest.event-history.enabled", false)
.containsEntry("lingniu.ingest.vehicle-state.enabled", false)

View File

@@ -48,9 +48,7 @@ lingniu:
store: ${SESSION_STORE:memory}
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}

View File

@@ -1,6 +1,7 @@
package com.lingniu.ingest.jt808app;
import com.lingniu.ingest.core.config.IngestCoreAutoConfiguration;
import com.lingniu.ingest.identity.MySqlVehicleIdentityService;
import com.lingniu.ingest.identity.config.VehicleIdentityAutoConfiguration;
import com.lingniu.ingest.protocol.jt808.codec.Jt808MessageDecoder;
import com.lingniu.ingest.protocol.jt808.config.Jt808AutoConfiguration;
@@ -37,6 +38,8 @@ class Jt808IngestAppCompositionTest {
.withPropertyValues(
"lingniu.ingest.jt808.enabled=true",
"lingniu.ingest.jt808.port=0",
"lingniu.ingest.identity.store=mysql",
"lingniu.ingest.identity.mysql.initialize-schema=false",
"lingniu.ingest.session.store=memory",
"lingniu.ingest.sink.mq.enabled=true",
"lingniu.ingest.sink.mq.type=kafka",
@@ -54,6 +57,7 @@ class Jt808IngestAppCompositionTest {
contextRunner.run(context -> {
assertThat(context).hasSingleBean(Jt808MessageDecoder.class);
assertThat(context).hasSingleBean(Jt808NettyServer.class);
assertThat(context).hasSingleBean(MySqlVehicleIdentityService.class);
assertThat(context).hasSingleBean(KafkaEventSink.class);
assertThat(context).hasSingleBean(KafkaEnvelopeDeadLetterSink.class);
assertThat(context).hasSingleBean(ArchiveStore.class);

View File

@@ -28,7 +28,7 @@ class Jt808IngestAppDefaultsTest {
.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.identity.store", "${VEHICLE_IDENTITY_STORE:file}")
.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.identity.mysql.jdbc-url",

View File

@@ -78,8 +78,6 @@ lingniu:
enabled: ${VEHICLE_STATE_ENABLED:false}
vehicle-stat:
enabled: ${VEHICLE_STAT_ENABLED:true}
repository-type: ${VEHICLE_STAT_REPOSITORY_TYPE:jdbc}
file-path: ${VEHICLE_STAT_FILE_PATH:./target/vehicle-stat/}
zone-id: ${VEHICLE_STAT_ZONE_ID:Asia/Shanghai}
jt808:
enabled: ${VEHICLE_STAT_JT808_MILEAGE_ENABLED:false}

View File

@@ -7,7 +7,7 @@ import com.lingniu.ingest.sink.mq.SinkMqAutoConfiguration;
import com.lingniu.ingest.vehiclestate.VehicleStateEnvelopeIngestor;
import com.lingniu.ingest.vehiclestate.config.VehicleStateAutoConfiguration;
import com.lingniu.ingest.vehiclestat.DailyVehicleStatService;
import com.lingniu.ingest.vehiclestat.FileVehicleStatRepository;
import com.lingniu.ingest.vehiclestat.JdbcVehicleStatMetricRepository;
import com.lingniu.ingest.vehiclestat.VehicleStatController;
import com.lingniu.ingest.vehiclestat.VehicleStatEnvelopeIngestor;
import com.lingniu.ingest.vehiclestat.VehicleStatEventProcessor;
@@ -16,22 +16,17 @@ import com.lingniu.ingest.vehiclestat.VehicleStatRepository;
import com.lingniu.ingest.vehiclestat.config.VehicleStatAutoConfiguration;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.springframework.boot.autoconfigure.AutoConfigurations;
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
import org.springframework.context.ApplicationContext;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.util.ClassUtils;
import java.nio.file.Path;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.Mockito.mock;
class VehicleAnalyticsAppCompositionTest {
@TempDir
Path tempDir;
@Test
void createsStatsBeansWithoutProtocolListenerOrEventFileStore() {
new ApplicationContextRunner()
@@ -41,6 +36,7 @@ class VehicleAnalyticsAppCompositionTest {
VehicleStatAutoConfiguration.class))
.withAllowBeanDefinitionOverriding(true)
.withBean("kafkaProducer", KafkaProducer.class, VehicleAnalyticsAppCompositionTest::kafkaProducer)
.withBean(JdbcTemplate.class, () -> mock(JdbcTemplate.class))
.withPropertyValues(
"lingniu.ingest.sink.mq.enabled=true",
"lingniu.ingest.sink.mq.type=kafka",
@@ -48,13 +44,12 @@ class VehicleAnalyticsAppCompositionTest {
"lingniu.ingest.sink.mq.consumer.enabled=false",
"lingniu.ingest.vehicle-state.enabled=false",
"lingniu.ingest.vehicle-stat.enabled=true",
"lingniu.ingest.vehicle-stat.file-path=" + tempDir.resolve("vehicle-stat"),
"lingniu.ingest.event-file-store.enabled=false",
"lingniu.ingest.event-history.enabled=false",
"lingniu.ingest.gb32960.enabled=false")
.run(context -> {
assertThat(context).hasSingleBean(VehicleStatRepository.class);
assertThat(context).hasSingleBean(FileVehicleStatRepository.class);
assertThat(context).hasSingleBean(JdbcVehicleStatMetricRepository.class);
assertThat(context).hasSingleBean(DailyVehicleStatService.class);
assertThat(context).hasSingleBean(VehicleStatEventProcessor.class);
assertThat(context).hasSingleBean(VehicleStatEnvelopeIngestor.class);

View File

@@ -25,7 +25,6 @@ class VehicleAnalyticsAppDefaultsTest {
.containsEntry("lingniu.ingest.event-file-store.enabled", false)
.containsEntry("lingniu.ingest.event-history.enabled", false)
.containsEntry("lingniu.ingest.vehicle-stat.enabled", "${VEHICLE_STAT_ENABLED:true}")
.containsEntry("lingniu.ingest.vehicle-stat.repository-type", "${VEHICLE_STAT_REPOSITORY_TYPE:jdbc}")
.containsEntry("lingniu.ingest.vehicle-state.enabled", "${VEHICLE_STATE_ENABLED:false}")
.containsEntry("lingniu.ingest.sink.mq.consumer.enabled", "${KAFKA_CONSUMER_ENABLED:true}")
.containsEntry(

View File

@@ -60,16 +60,13 @@ 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:
enabled: ${KAFKA_ENABLED:true}

View File

@@ -1,6 +1,7 @@
package com.lingniu.ingest.yutongmqttapp;
import com.lingniu.ingest.core.config.IngestCoreAutoConfiguration;
import com.lingniu.ingest.identity.MySqlVehicleIdentityService;
import com.lingniu.ingest.identity.config.VehicleIdentityAutoConfiguration;
import com.lingniu.ingest.inbound.mqtt.client.MqttEndpointManager;
import com.lingniu.ingest.inbound.mqtt.config.MqttInboundAutoConfiguration;
@@ -37,6 +38,8 @@ class YutongMqttAppCompositionTest {
"lingniu.ingest.mqtt.endpoints[0].uri=tcp://127.0.0.1:1883",
"lingniu.ingest.mqtt.endpoints[0].topic=/yutong/#",
"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",
@@ -54,6 +57,7 @@ class YutongMqttAppCompositionTest {
assertThat(context).hasSingleBean(MqttEndpointManager.class);
assertThat(context).hasSingleBean(MqttProfileRegistry.class);
assertThat(context).hasSingleBean(MqttRealtimeHandler.class);
assertThat(context).hasSingleBean(MySqlVehicleIdentityService.class);
assertThat(context).hasSingleBean(KafkaEventSink.class);
assertThat(context).hasSingleBean(KafkaEnvelopeDeadLetterSink.class);
assertThat(context).hasSingleBean(ArchiveStore.class);

View File

@@ -30,8 +30,9 @@ class YutongMqttAppDefaultsTest {
.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.identity.store", "${VEHICLE_IDENTITY_STORE:file}")
.containsEntry("lingniu.ingest.identity.mysql.table", "${VEHICLE_IDENTITY_MYSQL_TABLE:vehicle_identity_binding}")
.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)