feat: refresh mysql vehicle identity cache
All checks were successful
ci/woodpecker/push/woodpecker Pipeline was successful
All checks were successful
ci/woodpecker/push/woodpecker Pipeline was successful
This commit is contained in:
@@ -23,6 +23,7 @@ export TDENGINE_PASSWORD='<tdengine-password>'
|
|||||||
export VEHICLE_IDENTITY_MYSQL_JDBC_URL='jdbc:mysql://<mysql-host>:3306/vehicle_ingest'
|
export VEHICLE_IDENTITY_MYSQL_JDBC_URL='jdbc:mysql://<mysql-host>:3306/vehicle_ingest'
|
||||||
export VEHICLE_IDENTITY_MYSQL_USERNAME='<mysql-user>'
|
export VEHICLE_IDENTITY_MYSQL_USERNAME='<mysql-user>'
|
||||||
export VEHICLE_IDENTITY_MYSQL_PASSWORD='<mysql-password>'
|
export VEHICLE_IDENTITY_MYSQL_PASSWORD='<mysql-password>'
|
||||||
|
export VEHICLE_IDENTITY_MYSQL_REFRESH_INTERVAL='60s'
|
||||||
|
|
||||||
mkdir -p "$HOME/Library/LaunchAgents" "$PROJECT_ROOT/data"
|
mkdir -p "$HOME/Library/LaunchAgents" "$PROJECT_ROOT/data"
|
||||||
mkdir -p /tmp/lingniu-gb32960-live /tmp/lingniu-jt808-live /tmp/lingniu-history-live
|
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_JDBC_URL__#$VEHICLE_IDENTITY_MYSQL_JDBC_URL#g" \
|
||||||
-e "s#__VEHICLE_IDENTITY_MYSQL_USERNAME__#$VEHICLE_IDENTITY_MYSQL_USERNAME#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_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" \
|
"deploy/local/launchctl/com.lingniu.${service}.plist.template" \
|
||||||
> "$HOME/Library/LaunchAgents/com.lingniu.${service}.plist"
|
> "$HOME/Library/LaunchAgents/com.lingniu.${service}.plist"
|
||||||
done
|
done
|
||||||
|
|||||||
@@ -42,6 +42,8 @@
|
|||||||
<string>__VEHICLE_IDENTITY_MYSQL_USERNAME__</string>
|
<string>__VEHICLE_IDENTITY_MYSQL_USERNAME__</string>
|
||||||
<key>VEHICLE_IDENTITY_MYSQL_PASSWORD</key>
|
<key>VEHICLE_IDENTITY_MYSQL_PASSWORD</key>
|
||||||
<string>__VEHICLE_IDENTITY_MYSQL_PASSWORD__</string>
|
<string>__VEHICLE_IDENTITY_MYSQL_PASSWORD__</string>
|
||||||
|
<key>VEHICLE_IDENTITY_MYSQL_REFRESH_INTERVAL</key>
|
||||||
|
<string>__VEHICLE_IDENTITY_MYSQL_REFRESH_INTERVAL__</string>
|
||||||
<key>NACOS_CONFIG_ENABLED</key>
|
<key>NACOS_CONFIG_ENABLED</key>
|
||||||
<string>false</string>
|
<string>false</string>
|
||||||
<key>MANAGEMENT_HEALTH_REDIS_ENABLED</key>
|
<key>MANAGEMENT_HEALTH_REDIS_ENABLED</key>
|
||||||
|
|||||||
@@ -67,6 +67,7 @@ services:
|
|||||||
VEHICLE_IDENTITY_MYSQL_JDBC_URL: ${VEHICLE_IDENTITY_MYSQL_JDBC_URL:-}
|
VEHICLE_IDENTITY_MYSQL_JDBC_URL: ${VEHICLE_IDENTITY_MYSQL_JDBC_URL:-}
|
||||||
VEHICLE_IDENTITY_MYSQL_USERNAME: ${VEHICLE_IDENTITY_MYSQL_USERNAME:-}
|
VEHICLE_IDENTITY_MYSQL_USERNAME: ${VEHICLE_IDENTITY_MYSQL_USERNAME:-}
|
||||||
VEHICLE_IDENTITY_MYSQL_PASSWORD: ${VEHICLE_IDENTITY_MYSQL_PASSWORD:-}
|
VEHICLE_IDENTITY_MYSQL_PASSWORD: ${VEHICLE_IDENTITY_MYSQL_PASSWORD:-}
|
||||||
|
VEHICLE_IDENTITY_MYSQL_REFRESH_INTERVAL: ${VEHICLE_IDENTITY_MYSQL_REFRESH_INTERVAL:-60s}
|
||||||
ports:
|
ports:
|
||||||
- "${JT808_HTTP_PORT:-20400}:20400"
|
- "${JT808_HTTP_PORT:-20400}:20400"
|
||||||
- "${JT808_TCP_PORT:-808}:808"
|
- "${JT808_TCP_PORT:-808}:808"
|
||||||
|
|||||||
@@ -115,7 +115,7 @@ JT808 注册身份绑定:
|
|||||||
- 生产需要把 0x0100 注册信息维护到 MySQL 时,设置 `VEHICLE_IDENTITY_STORE=mysql`。
|
- 生产需要把 0x0100 注册信息维护到 MySQL 时,设置 `VEHICLE_IDENTITY_STORE=mysql`。
|
||||||
- 需要提供 `VEHICLE_IDENTITY_MYSQL_JDBC_URL`、`VEHICLE_IDENTITY_MYSQL_USERNAME`、`VEHICLE_IDENTITY_MYSQL_PASSWORD`。
|
- 需要提供 `VEHICLE_IDENTITY_MYSQL_JDBC_URL`、`VEHICLE_IDENTITY_MYSQL_USERNAME`、`VEHICLE_IDENTITY_MYSQL_PASSWORD`。
|
||||||
- 服务启动时自动创建 `vehicle_identity_binding` 表,按 `protocol + identifier_type + identifier_value` 维护 VIN 绑定。
|
- 服务启动时自动创建 `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。
|
- 绑定类型包括 `PHONE`、`DEVICE_ID`、`PLATE`;后续同一终端上报 raw/event 时会优先解析成已绑定 VIN。
|
||||||
|
|
||||||
健康检查:
|
健康检查:
|
||||||
|
|||||||
@@ -57,6 +57,7 @@ lingniu:
|
|||||||
username: ${VEHICLE_IDENTITY_MYSQL_USERNAME:}
|
username: ${VEHICLE_IDENTITY_MYSQL_USERNAME:}
|
||||||
password: ${VEHICLE_IDENTITY_MYSQL_PASSWORD:}
|
password: ${VEHICLE_IDENTITY_MYSQL_PASSWORD:}
|
||||||
driver-class-name: ${VEHICLE_IDENTITY_MYSQL_DRIVER_CLASS_NAME:com.mysql.cj.jdbc.Driver}
|
driver-class-name: ${VEHICLE_IDENTITY_MYSQL_DRIVER_CLASS_NAME:com.mysql.cj.jdbc.Driver}
|
||||||
|
refresh-interval: ${VEHICLE_IDENTITY_MYSQL_REFRESH_INTERVAL:60s}
|
||||||
sink:
|
sink:
|
||||||
mq:
|
mq:
|
||||||
enabled: ${KAFKA_ENABLED:true}
|
enabled: ${KAFKA_ENABLED:true}
|
||||||
|
|||||||
@@ -32,6 +32,7 @@ class Jt808IngestAppDefaultsTest {
|
|||||||
.containsEntry("lingniu.ingest.identity.mysql.table", "${VEHICLE_IDENTITY_MYSQL_TABLE:vehicle_identity_binding}")
|
.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.jdbc-url", "${VEHICLE_IDENTITY_MYSQL_JDBC_URL:}")
|
||||||
.containsEntry("lingniu.ingest.identity.mysql.username", "${VEHICLE_IDENTITY_MYSQL_USERNAME:}")
|
.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.enabled", "${SINK_ARCHIVE_ENABLED:true}")
|
||||||
.containsEntry("lingniu.ingest.sink.archive.path", "${SINK_ARCHIVE_PATH:./archive/}")
|
.containsEntry("lingniu.ingest.sink.archive.path", "${SINK_ARCHIVE_PATH:./archive/}")
|
||||||
.containsEntry("lingniu.ingest.event-file-store.enabled", false)
|
.containsEntry("lingniu.ingest.event-file-store.enabled", false)
|
||||||
|
|||||||
@@ -1,29 +1,43 @@
|
|||||||
package com.lingniu.ingest.identity;
|
package com.lingniu.ingest.identity;
|
||||||
|
|
||||||
import com.lingniu.ingest.api.ProtocolId;
|
import com.lingniu.ingest.api.ProtocolId;
|
||||||
|
import org.slf4j.Logger;
|
||||||
|
import org.slf4j.LoggerFactory;
|
||||||
|
|
||||||
import javax.sql.DataSource;
|
import javax.sql.DataSource;
|
||||||
|
import java.time.Duration;
|
||||||
import java.sql.Connection;
|
import java.sql.Connection;
|
||||||
import java.sql.PreparedStatement;
|
import java.sql.PreparedStatement;
|
||||||
import java.sql.ResultSet;
|
import java.sql.ResultSet;
|
||||||
import java.sql.SQLException;
|
import java.sql.SQLException;
|
||||||
import java.sql.Statement;
|
import java.sql.Statement;
|
||||||
import java.util.Locale;
|
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 DataSource dataSource;
|
||||||
private final String table;
|
private final String table;
|
||||||
private final InMemoryVehicleIdentityService index = new InMemoryVehicleIdentityService();
|
private final ScheduledExecutorService refresher;
|
||||||
|
private volatile InMemoryVehicleIdentityService index;
|
||||||
|
|
||||||
public MySqlVehicleIdentityService(DataSource dataSource, String table) {
|
public MySqlVehicleIdentityService(DataSource dataSource, String table) {
|
||||||
|
this(dataSource, table, Duration.ZERO);
|
||||||
|
}
|
||||||
|
|
||||||
|
public MySqlVehicleIdentityService(DataSource dataSource, String table, Duration refreshInterval) {
|
||||||
if (dataSource == null) {
|
if (dataSource == null) {
|
||||||
throw new IllegalArgumentException("dataSource must not be null");
|
throw new IllegalArgumentException("dataSource must not be null");
|
||||||
}
|
}
|
||||||
this.dataSource = dataSource;
|
this.dataSource = dataSource;
|
||||||
this.table = sanitizeTable(table);
|
this.table = sanitizeTable(table);
|
||||||
initializeSchema();
|
initializeSchema();
|
||||||
loadIndex();
|
this.index = loadIndex();
|
||||||
|
this.refresher = startRefresher(refreshInterval);
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
@@ -45,6 +59,17 @@ public final class MySqlVehicleIdentityService implements VehicleIdentityResolve
|
|||||||
return index.resolve(lookup);
|
return index.resolve(lookup);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public void refresh() {
|
||||||
|
index = loadIndex();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void close() {
|
||||||
|
if (refresher != null) {
|
||||||
|
refresher.shutdownNow();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
private void initializeSchema() {
|
private void initializeSchema() {
|
||||||
String sql = """
|
String sql = """
|
||||||
CREATE TABLE IF NOT EXISTS %s (
|
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;
|
String sql = "SELECT protocol, identifier_type, identifier_value, vin FROM " + table;
|
||||||
try (Connection connection = dataSource.getConnection();
|
try (Connection connection = dataSource.getConnection();
|
||||||
PreparedStatement statement = connection.prepareStatement(sql);
|
PreparedStatement statement = connection.prepareStatement(sql);
|
||||||
ResultSet resultSet = statement.executeQuery()) {
|
ResultSet resultSet = statement.executeQuery()) {
|
||||||
while (resultSet.next()) {
|
while (resultSet.next()) {
|
||||||
bindLoadedRow(
|
bindLoadedRow(
|
||||||
|
loaded,
|
||||||
protocol(resultSet.getString("protocol")),
|
protocol(resultSet.getString("protocol")),
|
||||||
resultSet.getString("identifier_type"),
|
resultSet.getString("identifier_type"),
|
||||||
resultSet.getString("identifier_value"),
|
resultSet.getString("identifier_value"),
|
||||||
resultSet.getString("vin"));
|
resultSet.getString("vin"));
|
||||||
}
|
}
|
||||||
|
return loaded;
|
||||||
} catch (SQLException e) {
|
} catch (SQLException e) {
|
||||||
throw new IllegalStateException("vehicle identity mysql index load failed: " + table, 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()) {
|
if (vin == null || vin.isBlank()) {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
@@ -95,7 +127,29 @@ public final class MySqlVehicleIdentityService implements VehicleIdentityResolve
|
|||||||
default -> null;
|
default -> null;
|
||||||
};
|
};
|
||||||
if (binding != 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);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -54,7 +54,8 @@ public class VehicleIdentityAutoConfiguration {
|
|||||||
@ConditionalOnProperty(prefix = "lingniu.ingest.identity", name = "store", havingValue = "mysql")
|
@ConditionalOnProperty(prefix = "lingniu.ingest.identity", name = "store", havingValue = "mysql")
|
||||||
public MySqlVehicleIdentityService mySqlVehicleIdentityService(VehicleIdentityProperties properties,
|
public MySqlVehicleIdentityService mySqlVehicleIdentityService(VehicleIdentityProperties properties,
|
||||||
DataSource dataSource) {
|
DataSource dataSource) {
|
||||||
return new MySqlVehicleIdentityService(dataSource, properties.getMysql().getTable());
|
VehicleIdentityProperties.Mysql mysql = properties.getMysql();
|
||||||
|
return new MySqlVehicleIdentityService(dataSource, mysql.getTable(), mysql.getRefreshInterval());
|
||||||
}
|
}
|
||||||
|
|
||||||
@Bean
|
@Bean
|
||||||
|
|||||||
@@ -2,6 +2,8 @@ package com.lingniu.ingest.identity.config;
|
|||||||
|
|
||||||
import org.springframework.boot.context.properties.ConfigurationProperties;
|
import org.springframework.boot.context.properties.ConfigurationProperties;
|
||||||
|
|
||||||
|
import java.time.Duration;
|
||||||
|
|
||||||
@ConfigurationProperties(prefix = "lingniu.ingest.identity")
|
@ConfigurationProperties(prefix = "lingniu.ingest.identity")
|
||||||
public class VehicleIdentityProperties {
|
public class VehicleIdentityProperties {
|
||||||
|
|
||||||
@@ -29,6 +31,7 @@ public class VehicleIdentityProperties {
|
|||||||
private String username = "";
|
private String username = "";
|
||||||
private String password = "";
|
private String password = "";
|
||||||
private String driverClassName = "com.mysql.cj.jdbc.Driver";
|
private String driverClassName = "com.mysql.cj.jdbc.Driver";
|
||||||
|
private Duration refreshInterval = Duration.ofSeconds(60);
|
||||||
|
|
||||||
public String getTable() { return table; }
|
public String getTable() { return table; }
|
||||||
public void setTable(String table) { this.table = 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 void setPassword(String password) { this.password = password; }
|
||||||
public String getDriverClassName() { return driverClassName; }
|
public String getDriverClassName() { return driverClassName; }
|
||||||
public void setDriverClassName(String driverClassName) { this.driverClassName = driverClassName; }
|
public void setDriverClassName(String driverClassName) { this.driverClassName = driverClassName; }
|
||||||
|
public Duration getRefreshInterval() { return refreshInterval; }
|
||||||
|
public void setRefreshInterval(Duration refreshInterval) { this.refreshInterval = refreshInterval; }
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -84,6 +84,27 @@ class MySqlVehicleIdentityServiceTest {
|
|||||||
assertThat(jdbc.lookupSelects).isZero();
|
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 static final class RecordingIdentityJdbc {
|
||||||
private final List<String> createdTables = new ArrayList<>();
|
private final List<String> createdTables = new ArrayList<>();
|
||||||
private final Map<String, String> bindings = new HashMap<>();
|
private final Map<String, String> bindings = new HashMap<>();
|
||||||
|
|||||||
Reference in New Issue
Block a user