diff --git a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/RawFrameHistoryController.java b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/RawFrameHistoryController.java index 4d48d07b..7f65e0c8 100644 --- a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/RawFrameHistoryController.java +++ b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/RawFrameHistoryController.java @@ -3,11 +3,13 @@ package com.lingniu.ingest.eventhistory; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.ObjectMapper; import com.lingniu.ingest.tdenginehistory.TdengineHistoryReader; +import com.lingniu.ingest.tdenginehistory.TdengineLocationRow; import com.lingniu.ingest.tdenginehistory.TdenginePage; import com.lingniu.ingest.tdenginehistory.TdenginePageCursor; import com.lingniu.ingest.tdenginehistory.TdengineQueryOrder; import com.lingniu.ingest.tdenginehistory.TdengineRawFrameQuery; import com.lingniu.ingest.tdenginehistory.TdengineRawFrameRow; +import com.lingniu.ingest.tdenginehistory.TdengineTelemetryFieldRow; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.Parameter; import io.swagger.v3.oas.annotations.tags.Tag; @@ -22,6 +24,7 @@ import org.springframework.web.server.ResponseStatusException; import java.io.IOException; import java.time.Instant; +import java.util.LinkedHashMap; import java.util.List; import java.util.Locale; import java.util.Map; @@ -35,6 +38,8 @@ public final class RawFrameHistoryController { private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); private static final TypeReference> MAP_TYPE = new TypeReference<>() {}; + private static final int LOCATION_ENRICH_LIMIT = 1; + private static final int TELEMETRY_ENRICH_LIMIT = 500; private final TdengineHistoryReader reader; @@ -71,6 +76,8 @@ public final class RawFrameHistoryController { @RequestParam(defaultValue = "DESC") TdengineQueryOrder order, @Parameter(description = "返回原始帧数量上限,最大 1000。", example = "100") @RequestParam(defaultValue = "100") int limit, + @Parameter(description = "是否按 rawUri 补充 location/telemetry 解析字段。默认 false,避免 raw 列表产生额外 TDengine 查询。", example = "false") + @RequestParam(defaultValue = "false") boolean includeParsedFields, @Parameter(description = "上一页返回的 nextCursor.ts。", example = "2026-06-29T05:00:01Z") @RequestParam(required = false) String cursorTs, @Parameter(description = "上一页返回的 nextCursor.id。", example = "frame-xxx") @@ -93,7 +100,7 @@ public final class RawFrameHistoryController { order, limit, cursor(cursorTs, cursorId))); - return RawFramePageResponse.from(page); + return RawFramePageResponse.from(reader, page, includeParsedFields); } private static Integer parseMessageId(String messageId) { @@ -144,12 +151,18 @@ public final class RawFrameHistoryController { List items, CursorResponse nextCursor) { - private static RawFramePageResponse from(TdenginePage page) { + private static RawFramePageResponse from(TdengineHistoryReader reader, + TdenginePage page, + boolean includeParsedFields) throws IOException { CursorResponse next = page.nextCursor() .map(cursor -> new CursorResponse(cursor.ts().toString(), cursor.tieBreaker())) .orElse(null); + List items = new java.util.ArrayList<>(page.items().size()); + for (TdengineRawFrameRow row : page.items()) { + items.add(RawFrameResponse.from(reader, row, includeParsedFields)); + } return new RawFramePageResponse( - page.items().stream().map(RawFrameResponse::from).toList(), + items, next); } } @@ -172,12 +185,16 @@ public final class RawFrameHistoryController { String peer, String metadataJson, Map metadata, + Map parsedFields, String protocol, String vehicleKey, String vin, String phone) { - private static RawFrameResponse from(TdengineRawFrameRow row) { + private static RawFrameResponse from(TdengineHistoryReader reader, + TdengineRawFrameRow row, + boolean includeParsedFields) throws IOException { + Map metadata = metadata(row.metadataJson()); return new RawFrameResponse( row.ts().toString(), row.receivedAt().toString(), @@ -192,7 +209,8 @@ public final class RawFrameHistoryController { row.parseError(), row.peer(), row.metadataJson(), - metadata(row.metadataJson()), + metadata, + parsedFields(reader, row, metadata, includeParsedFields), row.protocol(), row.vehicleKey(), row.vin(), @@ -209,5 +227,63 @@ public final class RawFrameHistoryController { return Map.of("_raw", metadataJson); } } + + private static Map parsedFields(TdengineHistoryReader reader, + TdengineRawFrameRow row, + Map metadata, + boolean includeParsedFields) throws IOException { + Map out = new LinkedHashMap<>(); + for (Map.Entry entry : metadata.entrySet()) { + String key = entry.getKey(); + if (key == null || key.isBlank()) { + continue; + } + if (key.startsWith("raw.") || key.startsWith("jt808.")) { + out.put(key, entry.getValue()); + } else { + out.put("raw." + key, entry.getValue()); + } + } + if (!includeParsedFields) { + return Map.copyOf(out); + } + for (TdengineLocationRow location : reader.queryLocationsByRawUri( + row.protocol(), row.rawUri(), LOCATION_ENRICH_LIMIT)) { + addLocation(out, location); + } + for (TdengineTelemetryFieldRow field : reader.queryTelemetryFieldsByRawUri( + row.protocol(), row.rawUri(), TELEMETRY_ENRICH_LIMIT)) { + out.put(field.fieldKey(), typedValue(field)); + out.put(field.fieldKey() + "._type", field.valueType()); + if (field.unit() != null && !field.unit().isBlank()) { + out.put(field.fieldKey() + "._unit", field.unit()); + } + } + return Map.copyOf(out); + } + + private static void addLocation(Map out, TdengineLocationRow row) { + out.put("location.longitude", row.longitude()); + out.put("location.latitude", row.latitude()); + out.put("location.altitudeM", row.altitudeM()); + out.put("location.speedKmh", row.speedKmh()); + out.put("location.directionDeg", row.directionDeg()); + out.put("location.alarmFlag", row.alarmFlag()); + out.put("location.statusFlag", row.statusFlag()); + if (row.totalMileageKm() != null) { + out.put("location.totalMileageKm", row.totalMileageKm()); + } + } + + private static Object typedValue(TdengineTelemetryFieldRow field) { + if (field.valueDouble() != null) { + return field.valueDouble(); + } + if (field.valueLong() != null) { + return field.valueLong(); + } + return field.valueText(); + } } + } diff --git a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/RawFrameHistoryControllerTest.java b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/RawFrameHistoryControllerTest.java index dbf40e44..eea46a31 100644 --- a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/RawFrameHistoryControllerTest.java +++ b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/RawFrameHistoryControllerTest.java @@ -16,6 +16,7 @@ import org.springframework.web.server.ResponseStatusException; import java.io.IOException; import java.time.Instant; +import java.util.ArrayList; import java.util.List; import java.util.Optional; @@ -42,6 +43,7 @@ class RawFrameHistoryControllerTest { "2026-06-29T23:59:59+08:00", TdengineQueryOrder.DESC, 20, + false, "2026-06-29T12:00:00Z", "frame-0"); @@ -61,10 +63,86 @@ class RawFrameHistoryControllerTest { .containsEntry("rawArchiveUri", item.rawUri()) .containsEntry("plateNo", "粤BD12345") .containsEntry("rawJson", "{\"carName\":\"粤BD12345\"}"); + assertThat(item.parsedFields()) + .containsEntry("raw.rawArchiveUri", item.rawUri()) + .containsEntry("raw.plateNo", "粤BD12345"); + assertThat(reader.locationRawUri).isNull(); + assertThat(reader.telemetryRawUri).isNull(); assertThat(response.nextCursor()).isEqualTo(new RawFrameHistoryController.CursorResponse( "2026-06-29T12:00:01Z", "frame-1")); } + @Test + void enrichesRawFrameWithLocationAndTelemetryFieldsByRawUri() throws Exception { + Instant ts = Instant.parse("2026-06-29T12:00:01Z"); + TdengineRawFrameRow raw = row("frame-1", ts.toString()); + CapturingReader reader = new CapturingReader(new TdenginePage<>(List.of(raw), Optional.empty())); + reader.locationsByRawUri.add(new TdengineLocationRow( + ts, + "loc-1", + "frame-1", + ts.plusMillis(500), + 121.312516, + 30.832808, + 6.0, + 25.0, + 333.0, + 0, + 524290, + 14504.446, + raw.rawUri(), + "{\"rawArchiveUri\":\"" + raw.rawUri() + "\"}", + "JT808", + "VIN-JT808-1", + "VIN-JT808-1", + "g7gps")); + reader.telemetryByRawUri.add(new TdengineTelemetryFieldRow( + ts, + "tf-1", + "frame-1", + ts.plusMillis(500), + "VEHICLE.speedKmh", + "DOUBLE", + "25.0", + 25.0, + null, + "km/h", + "GOOD", + "VEHICLE.speedKmh", + raw.rawUri(), + "{}", + "GB32960", + "VIN-GB-1", + "VIN-GB-1", + "")); + RawFrameHistoryController controller = new RawFrameHistoryController(reader); + + RawFrameHistoryController.RawFramePageResponse response = controller.rawFrames( + "jt808", + null, + "VIN-JT808-1", + null, + "0x0200", + null, + "2026-06-29T00:00:00+08:00", + "2026-06-29T23:59:59+08:00", + TdengineQueryOrder.DESC, + 20, + true, + null, + null); + + RawFrameHistoryController.RawFrameResponse item = response.items().getFirst(); + assertThat(reader.locationRawUri).isEqualTo(raw.rawUri()); + assertThat(reader.telemetryRawUri).isEqualTo(raw.rawUri()); + assertThat(item.parsedFields()) + .containsEntry("location.longitude", 121.312516) + .containsEntry("location.latitude", 30.832808) + .containsEntry("location.speedKmh", 25.0) + .containsEntry("location.totalMileageKm", 14504.446) + .containsEntry("VEHICLE.speedKmh", 25.0); + } + @Test void rejectsRawFrameQueryWithoutAnyNarrowingFilter() { RawFrameHistoryController controller = new RawFrameHistoryController( @@ -81,6 +159,7 @@ class RawFrameHistoryControllerTest { "2026-06-30", TdengineQueryOrder.DESC, 20, + false, null, null)) .isInstanceOfSatisfying(ResponseStatusException.class, ex -> @@ -112,7 +191,11 @@ class RawFrameHistoryControllerTest { private static final class CapturingReader implements TdengineHistoryReader { private final TdenginePage response; + private final List locationsByRawUri = new ArrayList<>(); + private final List telemetryByRawUri = new ArrayList<>(); private TdengineRawFrameQuery query; + private String locationRawUri; + private String telemetryRawUri; private CapturingReader(TdenginePage response) { this.response = response; @@ -133,5 +216,17 @@ class RawFrameHistoryControllerTest { public TdenginePage queryTelemetryFields(TdengineTelemetryFieldQuery query) { throw new UnsupportedOperationException(); } + + @Override + public List queryLocationsByRawUri(String protocol, String rawUri, int limit) { + locationRawUri = rawUri; + return locationsByRawUri; + } + + @Override + public List queryTelemetryFieldsByRawUri(String protocol, String rawUri, int limit) { + telemetryRawUri = rawUri; + return telemetryByRawUri; + } } } diff --git a/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineHistoryQueries.java b/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineHistoryQueries.java index 3555119c..f2a4d756 100644 --- a/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineHistoryQueries.java +++ b/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineHistoryQueries.java @@ -45,6 +45,14 @@ public final class TdengineHistoryQueries { query.order(), query.limit(), query.cursor()); } + public TdengineQueryStatement locationsByRawUri(String protocol, String rawUri, int limit) { + return byRawUri(schema.locationsStableTable(), LOCATION_COLUMNS, protocol, rawUri, "fact_id", limit); + } + + public TdengineQueryStatement telemetryFieldsByRawUri(String protocol, String rawUri, int limit) { + return byRawUri(schema.telemetryFieldsStableTable(), TELEMETRY_FIELD_COLUMNS, protocol, rawUri, "fact_id", limit); + } + private static TdengineQueryStatement query(String table, String columns, String tieBreakerColumn, @@ -113,4 +121,22 @@ public final class TdengineHistoryQueries { values.add(value); } } + + private static TdengineQueryStatement byRawUri(String table, + String columns, + String protocol, + String rawUri, + String tieBreakerColumn, + int limit) { + List values = new ArrayList<>(); + StringBuilder sql = new StringBuilder(256) + .append("SELECT ").append(columns) + .append(" FROM ").append(table) + .append(" WHERE raw_uri = ?"); + values.add(rawUri); + addFilter(sql, values, "protocol", protocol == null ? "" : protocol.trim().toUpperCase()); + sql.append(" ORDER BY ts DESC, ").append(tieBreakerColumn).append(" DESC LIMIT ?"); + values.add(Math.max(1, Math.min(limit, 1000))); + return new TdengineQueryStatement(sql.toString(), values); + } } diff --git a/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineHistoryReader.java b/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineHistoryReader.java index ec050748..d12c8824 100644 --- a/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineHistoryReader.java +++ b/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineHistoryReader.java @@ -1,6 +1,7 @@ package com.lingniu.ingest.tdenginehistory; import java.io.IOException; +import java.util.List; public interface TdengineHistoryReader { @@ -9,4 +10,14 @@ public interface TdengineHistoryReader { TdenginePage queryLocations(TdengineLocationQuery query) throws IOException; TdenginePage queryTelemetryFields(TdengineTelemetryFieldQuery query) throws IOException; + + default List queryLocationsByRawUri(String protocol, String rawUri, int limit) + throws IOException { + return List.of(); + } + + default List queryTelemetryFieldsByRawUri(String protocol, String rawUri, int limit) + throws IOException { + return List.of(); + } } diff --git a/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineHistorySchema.java b/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineHistorySchema.java index b51d0829..66714ed0 100644 --- a/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineHistorySchema.java +++ b/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineHistorySchema.java @@ -31,6 +31,14 @@ public final class TdengineHistorySchema { return "raw_frames"; } + public String locationsStableTable() { + return "vehicle_locations"; + } + + public String telemetryFieldsStableTable() { + return "telemetry_fields"; + } + public String locationTable(String protocol, String vehicleKey) { return "loc_" + TdengineIdentifier.fragment(protocol) + "_" + TdengineIdentifier.hash16(vehicleKey); } diff --git a/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineJdbcHistoryReader.java b/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineJdbcHistoryReader.java index 09e09dd9..8779971a 100644 --- a/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineJdbcHistoryReader.java +++ b/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineJdbcHistoryReader.java @@ -49,6 +49,26 @@ public final class TdengineJdbcHistoryReader implements TdengineHistoryReader { row -> new TdenginePageCursor(row.ts(), row.factId())); } + @Override + public List queryLocationsByRawUri(String protocol, String rawUri, int limit) + throws IOException { + if (rawUri == null || rawUri.isBlank()) { + return List.of(); + } + TdengineQueryStatement statement = queries.locationsByRawUri(protocol, rawUri, limit); + return readList(statement, this::location); + } + + @Override + public List queryTelemetryFieldsByRawUri(String protocol, String rawUri, int limit) + throws IOException { + if (rawUri == null || rawUri.isBlank()) { + return List.of(); + } + TdengineQueryStatement statement = queries.telemetryFieldsByRawUri(protocol, rawUri, limit); + return readList(statement, this::telemetryField); + } + private TdenginePage readPage(TdengineQueryStatement statement, int pageSize, SqlRowMapper mapper, @@ -76,6 +96,25 @@ public final class TdengineJdbcHistoryReader implements TdengineHistoryReader { return new TdenginePage<>(items, next); } + private List readList(TdengineQueryStatement statement, SqlRowMapper mapper) throws IOException { + List rows = new ArrayList<>(); + try (Connection connection = dataSource.getConnection(); + PreparedStatement prepared = connection.prepareStatement(statement.sql())) { + bind(prepared, statement.values()); + try (ResultSet resultSet = prepared.executeQuery()) { + while (resultSet.next()) { + rows.add(mapper.map(resultSet)); + } + } + } catch (SQLException e) { + if (isTableMissing(e)) { + return List.of(); + } + throw new IOException("tdengine history query failed", e); + } + return List.copyOf(rows); + } + private static boolean isTableMissing(SQLException e) { for (Throwable current = e; current != null; current = current.getCause()) { String message = current.getMessage();