diff --git a/deploy/local/launchctl/README.md b/deploy/local/launchctl/README.md index b15caf73..51d68cad 100644 --- a/deploy/local/launchctl/README.md +++ b/deploy/local/launchctl/README.md @@ -23,6 +23,7 @@ export TDENGINE_PASSWORD='' export VEHICLE_IDENTITY_MYSQL_JDBC_URL='jdbc:mysql://:3306/vehicle_ingest' export VEHICLE_IDENTITY_MYSQL_USERNAME='' export VEHICLE_IDENTITY_MYSQL_PASSWORD='' +export VEHICLE_IDENTITY_MYSQL_REFRESH_INTERVAL='60s' mkdir -p "$HOME/Library/LaunchAgents" "$PROJECT_ROOT/data" mkdir -p /tmp/lingniu-gb32960-live /tmp/lingniu-jt808-live /tmp/lingniu-history-live @@ -39,6 +40,7 @@ for service in gb32960 jt808 vehicle-history; do -e "s#__VEHICLE_IDENTITY_MYSQL_JDBC_URL__#$VEHICLE_IDENTITY_MYSQL_JDBC_URL#g" \ -e "s#__VEHICLE_IDENTITY_MYSQL_USERNAME__#$VEHICLE_IDENTITY_MYSQL_USERNAME#g" \ -e "s#__VEHICLE_IDENTITY_MYSQL_PASSWORD__#$VEHICLE_IDENTITY_MYSQL_PASSWORD#g" \ + -e "s#__VEHICLE_IDENTITY_MYSQL_REFRESH_INTERVAL__#$VEHICLE_IDENTITY_MYSQL_REFRESH_INTERVAL#g" \ "deploy/local/launchctl/com.lingniu.${service}.plist.template" \ > "$HOME/Library/LaunchAgents/com.lingniu.${service}.plist" done diff --git a/deploy/local/launchctl/com.lingniu.jt808.plist.template b/deploy/local/launchctl/com.lingniu.jt808.plist.template index 731bc532..4773a3f3 100644 --- a/deploy/local/launchctl/com.lingniu.jt808.plist.template +++ b/deploy/local/launchctl/com.lingniu.jt808.plist.template @@ -42,6 +42,8 @@ __VEHICLE_IDENTITY_MYSQL_USERNAME__ VEHICLE_IDENTITY_MYSQL_PASSWORD __VEHICLE_IDENTITY_MYSQL_PASSWORD__ + VEHICLE_IDENTITY_MYSQL_REFRESH_INTERVAL + __VEHICLE_IDENTITY_MYSQL_REFRESH_INTERVAL__ NACOS_CONFIG_ENABLED false MANAGEMENT_HEALTH_REDIS_ENABLED diff --git a/deploy/portainer/docker-compose.yml b/deploy/portainer/docker-compose.yml index 34f14275..ece55207 100644 --- a/deploy/portainer/docker-compose.yml +++ b/deploy/portainer/docker-compose.yml @@ -67,6 +67,7 @@ services: VEHICLE_IDENTITY_MYSQL_JDBC_URL: ${VEHICLE_IDENTITY_MYSQL_JDBC_URL:-} VEHICLE_IDENTITY_MYSQL_USERNAME: ${VEHICLE_IDENTITY_MYSQL_USERNAME:-} VEHICLE_IDENTITY_MYSQL_PASSWORD: ${VEHICLE_IDENTITY_MYSQL_PASSWORD:-} + VEHICLE_IDENTITY_MYSQL_REFRESH_INTERVAL: ${VEHICLE_IDENTITY_MYSQL_REFRESH_INTERVAL:-60s} ports: - "${JT808_HTTP_PORT:-20400}:20400" - "${JT808_TCP_PORT:-808}:808" diff --git a/docs/operations/vehicle-ingest-tdengine-verification.md b/docs/operations/vehicle-ingest-tdengine-verification.md index d5b18c1d..66613579 100644 --- a/docs/operations/vehicle-ingest-tdengine-verification.md +++ b/docs/operations/vehicle-ingest-tdengine-verification.md @@ -115,7 +115,7 @@ JT808 注册身份绑定: - 生产需要把 0x0100 注册信息维护到 MySQL 时,设置 `VEHICLE_IDENTITY_STORE=mysql`。 - 需要提供 `VEHICLE_IDENTITY_MYSQL_JDBC_URL`、`VEHICLE_IDENTITY_MYSQL_USERNAME`、`VEHICLE_IDENTITY_MYSQL_PASSWORD`。 - 服务启动时自动创建 `vehicle_identity_binding` 表,按 `protocol + identifier_type + identifier_value` 维护 VIN 绑定。 -- MySQL 绑定在启动时加载到内存索引,注册/补写时写穿 MySQL 并刷新内存索引;帧解析热路径只查内存,不按帧访问 MySQL。 +- MySQL 绑定在启动时加载到内存索引,注册/补写时写穿 MySQL 并刷新内存索引;帧解析热路径只查内存,不按帧访问 MySQL。外部系统更新绑定后,后台按 `VEHICLE_IDENTITY_MYSQL_REFRESH_INTERVAL` 周期刷新,默认 `60s`。 - 绑定类型包括 `PHONE`、`DEVICE_ID`、`PLATE`;后续同一终端上报 raw/event 时会优先解析成已绑定 VIN。 健康检查: 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 cb83c7aa..88606230 100644 --- a/modules/apps/jt808-ingest-app/src/main/resources/application.yml +++ b/modules/apps/jt808-ingest-app/src/main/resources/application.yml @@ -57,6 +57,7 @@ lingniu: username: ${VEHICLE_IDENTITY_MYSQL_USERNAME:} 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} sink: mq: enabled: ${KAFKA_ENABLED:true} 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 281c54b2..458b4258 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 @@ -32,6 +32,7 @@ class Jt808IngestAppDefaultsTest { .containsEntry("lingniu.ingest.identity.mysql.table", "${VEHICLE_IDENTITY_MYSQL_TABLE:vehicle_identity_binding}") .containsEntry("lingniu.ingest.identity.mysql.jdbc-url", "${VEHICLE_IDENTITY_MYSQL_JDBC_URL:}") .containsEntry("lingniu.ingest.identity.mysql.username", "${VEHICLE_IDENTITY_MYSQL_USERNAME:}") + .containsEntry("lingniu.ingest.identity.mysql.refresh-interval", "${VEHICLE_IDENTITY_MYSQL_REFRESH_INTERVAL:60s}") .containsEntry("lingniu.ingest.sink.archive.enabled", "${SINK_ARCHIVE_ENABLED:true}") .containsEntry("lingniu.ingest.sink.archive.path", "${SINK_ARCHIVE_PATH:./archive/}") .containsEntry("lingniu.ingest.event-file-store.enabled", false) diff --git a/modules/core/vehicle-identity/src/main/java/com/lingniu/ingest/identity/MySqlVehicleIdentityService.java b/modules/core/vehicle-identity/src/main/java/com/lingniu/ingest/identity/MySqlVehicleIdentityService.java index 634819d2..4f31745e 100644 --- a/modules/core/vehicle-identity/src/main/java/com/lingniu/ingest/identity/MySqlVehicleIdentityService.java +++ b/modules/core/vehicle-identity/src/main/java/com/lingniu/ingest/identity/MySqlVehicleIdentityService.java @@ -1,29 +1,43 @@ package com.lingniu.ingest.identity; import com.lingniu.ingest.api.ProtocolId; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import javax.sql.DataSource; +import java.time.Duration; import java.sql.Connection; import java.sql.PreparedStatement; import java.sql.ResultSet; import java.sql.SQLException; import java.sql.Statement; import java.util.Locale; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; -public final class MySqlVehicleIdentityService implements VehicleIdentityResolver, VehicleIdentityRegistry { +public final class MySqlVehicleIdentityService implements VehicleIdentityResolver, VehicleIdentityRegistry, AutoCloseable { + + private static final Logger log = LoggerFactory.getLogger(MySqlVehicleIdentityService.class); private final DataSource dataSource; private final String table; - private final InMemoryVehicleIdentityService index = new InMemoryVehicleIdentityService(); + private final ScheduledExecutorService refresher; + private volatile InMemoryVehicleIdentityService index; public MySqlVehicleIdentityService(DataSource dataSource, String table) { + this(dataSource, table, Duration.ZERO); + } + + public MySqlVehicleIdentityService(DataSource dataSource, String table, Duration refreshInterval) { if (dataSource == null) { throw new IllegalArgumentException("dataSource must not be null"); } this.dataSource = dataSource; this.table = sanitizeTable(table); initializeSchema(); - loadIndex(); + this.index = loadIndex(); + this.refresher = startRefresher(refreshInterval); } @Override @@ -45,6 +59,17 @@ public final class MySqlVehicleIdentityService implements VehicleIdentityResolve return index.resolve(lookup); } + public void refresh() { + index = loadIndex(); + } + + @Override + public void close() { + if (refresher != null) { + refresher.shutdownNow(); + } + } + private void initializeSchema() { String sql = """ CREATE TABLE IF NOT EXISTS %s ( @@ -66,24 +91,31 @@ public final class MySqlVehicleIdentityService implements VehicleIdentityResolve } } - private void loadIndex() { + private InMemoryVehicleIdentityService loadIndex() { + InMemoryVehicleIdentityService loaded = new InMemoryVehicleIdentityService(); String sql = "SELECT protocol, identifier_type, identifier_value, vin FROM " + table; try (Connection connection = dataSource.getConnection(); PreparedStatement statement = connection.prepareStatement(sql); ResultSet resultSet = statement.executeQuery()) { while (resultSet.next()) { bindLoadedRow( + loaded, protocol(resultSet.getString("protocol")), resultSet.getString("identifier_type"), resultSet.getString("identifier_value"), resultSet.getString("vin")); } + return loaded; } catch (SQLException e) { throw new IllegalStateException("vehicle identity mysql index load failed: " + table, e); } } - private void bindLoadedRow(ProtocolId protocol, String identifierType, String identifierValue, String vin) { + private void bindLoadedRow(InMemoryVehicleIdentityService loaded, + ProtocolId protocol, + String identifierType, + String identifierValue, + String vin) { if (vin == null || vin.isBlank()) { return; } @@ -95,7 +127,29 @@ public final class MySqlVehicleIdentityService implements VehicleIdentityResolve default -> null; }; if (binding != null) { - index.bind(binding); + loaded.bind(binding); + } + } + + private ScheduledExecutorService startRefresher(Duration refreshInterval) { + if (refreshInterval == null || refreshInterval.isZero() || refreshInterval.isNegative()) { + return null; + } + long delayMillis = Math.max(refreshInterval.toMillis(), 1000L); + ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(r -> { + Thread t = new Thread(r, "vehicle-identity-mysql-refresh"); + t.setDaemon(true); + return t; + }); + executor.scheduleWithFixedDelay(this::safeRefresh, delayMillis, delayMillis, TimeUnit.MILLISECONDS); + return executor; + } + + private void safeRefresh() { + try { + refresh(); + } catch (RuntimeException e) { + log.warn("vehicle identity mysql refresh failed table={}", table, e); } } diff --git a/modules/core/vehicle-identity/src/main/java/com/lingniu/ingest/identity/config/VehicleIdentityAutoConfiguration.java b/modules/core/vehicle-identity/src/main/java/com/lingniu/ingest/identity/config/VehicleIdentityAutoConfiguration.java index fb5d858c..87407ffb 100644 --- a/modules/core/vehicle-identity/src/main/java/com/lingniu/ingest/identity/config/VehicleIdentityAutoConfiguration.java +++ b/modules/core/vehicle-identity/src/main/java/com/lingniu/ingest/identity/config/VehicleIdentityAutoConfiguration.java @@ -54,7 +54,8 @@ public class VehicleIdentityAutoConfiguration { @ConditionalOnProperty(prefix = "lingniu.ingest.identity", name = "store", havingValue = "mysql") public MySqlVehicleIdentityService mySqlVehicleIdentityService(VehicleIdentityProperties properties, DataSource dataSource) { - return new MySqlVehicleIdentityService(dataSource, properties.getMysql().getTable()); + VehicleIdentityProperties.Mysql mysql = properties.getMysql(); + return new MySqlVehicleIdentityService(dataSource, mysql.getTable(), mysql.getRefreshInterval()); } @Bean diff --git a/modules/core/vehicle-identity/src/main/java/com/lingniu/ingest/identity/config/VehicleIdentityProperties.java b/modules/core/vehicle-identity/src/main/java/com/lingniu/ingest/identity/config/VehicleIdentityProperties.java index 146efe74..9b1a98d4 100644 --- a/modules/core/vehicle-identity/src/main/java/com/lingniu/ingest/identity/config/VehicleIdentityProperties.java +++ b/modules/core/vehicle-identity/src/main/java/com/lingniu/ingest/identity/config/VehicleIdentityProperties.java @@ -2,6 +2,8 @@ package com.lingniu.ingest.identity.config; import org.springframework.boot.context.properties.ConfigurationProperties; +import java.time.Duration; + @ConfigurationProperties(prefix = "lingniu.ingest.identity") public class VehicleIdentityProperties { @@ -29,6 +31,7 @@ public class VehicleIdentityProperties { private String username = ""; private String password = ""; private String driverClassName = "com.mysql.cj.jdbc.Driver"; + private Duration refreshInterval = Duration.ofSeconds(60); public String getTable() { return table; } public void setTable(String table) { this.table = table; } @@ -40,5 +43,7 @@ public class VehicleIdentityProperties { public void setPassword(String password) { this.password = password; } public String getDriverClassName() { return driverClassName; } public void setDriverClassName(String driverClassName) { this.driverClassName = driverClassName; } + public Duration getRefreshInterval() { return refreshInterval; } + public void setRefreshInterval(Duration refreshInterval) { this.refreshInterval = refreshInterval; } } } diff --git a/modules/core/vehicle-identity/src/test/java/com/lingniu/ingest/identity/MySqlVehicleIdentityServiceTest.java b/modules/core/vehicle-identity/src/test/java/com/lingniu/ingest/identity/MySqlVehicleIdentityServiceTest.java index f4326f47..c4bc45ea 100644 --- a/modules/core/vehicle-identity/src/test/java/com/lingniu/ingest/identity/MySqlVehicleIdentityServiceTest.java +++ b/modules/core/vehicle-identity/src/test/java/com/lingniu/ingest/identity/MySqlVehicleIdentityServiceTest.java @@ -84,6 +84,27 @@ class MySqlVehicleIdentityServiceTest { assertThat(jdbc.lookupSelects).isZero(); } + @Test + void refreshLoadsExternallyInsertedBindingWithoutPerFrameDatabaseLookup() { + RecordingIdentityJdbc jdbc = new RecordingIdentityJdbc(); + MySqlVehicleIdentityService service = new MySqlVehicleIdentityService(jdbc.dataSource(), "vehicle_identity_binding"); + + VehicleIdentity before = service.resolve(new VehicleIdentityLookup( + ProtocolId.JT808, "", "13079960002", "", "")); + + jdbc.seed("JT808", "PHONE", "13079960002", "LNVIN000000000456"); + service.refresh(); + VehicleIdentity after = service.resolve(new VehicleIdentityLookup( + ProtocolId.JT808, "", "13079960002", "", "")); + + assertThat(before.resolved()).isFalse(); + assertThat(after.vin()).isEqualTo("LNVIN000000000456"); + assertThat(after.resolved()).isTrue(); + assertThat(after.source()).isEqualTo(VehicleIdentitySource.BOUND_PHONE); + assertThat(jdbc.loadAllSelects).isEqualTo(2); + assertThat(jdbc.lookupSelects).isZero(); + } + private static final class RecordingIdentityJdbc { private final List createdTables = new ArrayList<>(); private final Map bindings = new HashMap<>();