From 90342fe5bf0c00e16764b3f3f1db96274558751c Mon Sep 17 00:00:00 2001 From: lingniu Date: Mon, 29 Jun 2026 13:55:25 +0800 Subject: [PATCH] feat: query telemetry fields from tdengine --- .../Jt808LocationHistoryControllerTest.java | 6 +++ .../TdengineHistoryQueries.java | 9 ++++ .../TdengineHistoryReader.java | 2 + .../TdengineJdbcHistoryReader.java | 37 ++++++++++++++ .../TdengineTelemetryFieldQuery.java | 35 +++++++++++++ .../TdengineHistoryQueriesTest.java | 27 ++++++++++ .../TdengineJdbcHistoryReaderTest.java | 50 +++++++++++++++++++ 7 files changed, 166 insertions(+) create mode 100644 modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineTelemetryFieldQuery.java diff --git a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/Jt808LocationHistoryControllerTest.java b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/Jt808LocationHistoryControllerTest.java index d7624077..277862f4 100644 --- a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/Jt808LocationHistoryControllerTest.java +++ b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/Jt808LocationHistoryControllerTest.java @@ -109,5 +109,11 @@ class Jt808LocationHistoryControllerTest { this.query = query; return response; } + + @Override + public TdenginePage queryTelemetryFields( + com.lingniu.ingest.tdenginehistory.TdengineTelemetryFieldQuery query) { + throw new UnsupportedOperationException(); + } } } 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 91fcd539..d9b58694 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 @@ -11,6 +11,9 @@ public final class TdengineHistoryQueries { private static final String LOCATION_COLUMNS = "ts, fact_id, frame_id, received_at, longitude, latitude, " + "altitude_m, speed_kmh, direction_deg, alarm_flag, status_flag, total_mileage_km, raw_uri, " + "metadata_json, protocol, vehicle_key, vin, phone"; + private static final String TELEMETRY_FIELD_COLUMNS = "ts, fact_id, frame_id, received_at, field_key, " + + "value_type, value_text, value_double, value_long, unit, quality, source_path, raw_uri, " + + "metadata_json, protocol, vehicle_key, vin, phone"; private final TdengineHistorySchema schema; @@ -33,6 +36,12 @@ public final class TdengineHistoryQueries { query.order(), query.limit(), query.cursor()); } + public TdengineQueryStatement telemetryFields(TdengineTelemetryFieldQuery query) { + String table = schema.telemetryFieldTable(query.protocol(), query.vehicleKey(), query.fieldKey()); + return query(table, TELEMETRY_FIELD_COLUMNS, "fact_id", query.from(), query.to(), + query.order(), query.limit(), query.cursor()); + } + private static TdengineQueryStatement query(String table, String columns, String tieBreakerColumn, 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 bd99fe11..ec050748 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 @@ -7,4 +7,6 @@ public interface TdengineHistoryReader { TdenginePage queryRawFrames(TdengineRawFrameQuery query) throws IOException; TdenginePage queryLocations(TdengineLocationQuery query) throws IOException; + + TdenginePage queryTelemetryFields(TdengineTelemetryFieldQuery query) throws IOException; } 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 974be5ca..c6e16d14 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 @@ -40,6 +40,15 @@ public final class TdengineJdbcHistoryReader implements TdengineHistoryReader { return readPage(statement, pageSize, this::location, row -> new TdenginePageCursor(row.ts(), row.factId())); } + @Override + public TdenginePage queryTelemetryFields(TdengineTelemetryFieldQuery query) + throws IOException { + int pageSize = Math.min(query.limit(), 1000); + TdengineQueryStatement statement = queries.telemetryFields(query.withLimit(pageSize + 1)); + return readPage(statement, pageSize, this::telemetryField, + row -> new TdenginePageCursor(row.ts(), row.factId())); + } + private TdenginePage readPage(TdengineQueryStatement statement, int pageSize, SqlRowMapper mapper, @@ -121,6 +130,29 @@ public final class TdengineJdbcHistoryReader implements TdengineHistoryReader { ); } + private TdengineTelemetryFieldRow telemetryField(ResultSet rs) throws SQLException { + return new TdengineTelemetryFieldRow( + instant(rs, "ts"), + rs.getString("fact_id"), + rs.getString("frame_id"), + instant(rs, "received_at"), + rs.getString("field_key"), + rs.getString("value_type"), + rs.getString("value_text"), + nullableDouble(rs, "value_double"), + nullableLong(rs, "value_long"), + rs.getString("unit"), + rs.getString("quality"), + rs.getString("source_path"), + rs.getString("raw_uri"), + rs.getString("metadata_json"), + rs.getString("protocol"), + rs.getString("vehicle_key"), + rs.getString("vin"), + rs.getString("phone") + ); + } + private static Instant instant(ResultSet rs, String column) throws SQLException { Timestamp timestamp = rs.getTimestamp(column); return timestamp == null ? null : timestamp.toInstant(); @@ -131,6 +163,11 @@ public final class TdengineJdbcHistoryReader implements TdengineHistoryReader { return value instanceof Number number ? number.doubleValue() : null; } + private static Long nullableLong(ResultSet rs, String column) throws SQLException { + Object value = rs.getObject(column); + return value instanceof Number number ? number.longValue() : null; + } + @FunctionalInterface private interface SqlRowMapper { T map(ResultSet rs) throws SQLException; diff --git a/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineTelemetryFieldQuery.java b/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineTelemetryFieldQuery.java new file mode 100644 index 00000000..dcdeb8ab --- /dev/null +++ b/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineTelemetryFieldQuery.java @@ -0,0 +1,35 @@ +package com.lingniu.ingest.tdenginehistory; + +import java.time.Instant; + +public record TdengineTelemetryFieldQuery( + String protocol, + String vehicleKey, + String fieldKey, + Instant from, + Instant to, + TdengineQueryOrder order, + int limit, + TdenginePageCursor cursor +) { + public TdengineTelemetryFieldQuery { + if (protocol == null || protocol.isBlank()) { + throw new IllegalArgumentException("protocol must not be blank"); + } + if (vehicleKey == null || vehicleKey.isBlank()) { + throw new IllegalArgumentException("vehicleKey must not be blank"); + } + if (fieldKey == null || fieldKey.isBlank()) { + throw new IllegalArgumentException("fieldKey must not be blank"); + } + if (from == null || to == null || !from.isBefore(to)) { + throw new IllegalArgumentException("query time range must be valid"); + } + order = order == null ? TdengineQueryOrder.DESC : order; + limit = Math.max(1, Math.min(limit, 1001)); + } + + TdengineTelemetryFieldQuery withLimit(int newLimit) { + return new TdengineTelemetryFieldQuery(protocol, vehicleKey, fieldKey, from, to, order, newLimit, cursor); + } +} diff --git a/modules/sinks/tdengine-history-store/src/test/java/com/lingniu/ingest/tdenginehistory/TdengineHistoryQueriesTest.java b/modules/sinks/tdengine-history-store/src/test/java/com/lingniu/ingest/tdenginehistory/TdengineHistoryQueriesTest.java index 918d0364..ac04677d 100644 --- a/modules/sinks/tdengine-history-store/src/test/java/com/lingniu/ingest/tdenginehistory/TdengineHistoryQueriesTest.java +++ b/modules/sinks/tdengine-history-store/src/test/java/com/lingniu/ingest/tdenginehistory/TdengineHistoryQueriesTest.java @@ -56,4 +56,31 @@ class TdengineHistoryQueriesTest { .contains(" ORDER BY ts DESC, fact_id DESC LIMIT ?"); assertThat(statement.values().getLast()).isEqualTo(1001); } + + @Test + void telemetryFieldQueryUsesFieldChildTableAndKeysetPagination() { + TdengineTelemetryFieldQuery query = new TdengineTelemetryFieldQuery( + "GB32960", + "VIN123", + "VEHICLE.totalMileageKm", + Instant.parse("2026-06-29T00:00:00Z"), + Instant.parse("2026-06-30T00:00:00Z"), + TdengineQueryOrder.ASC, + 100, + new TdenginePageCursor(Instant.parse("2026-06-29T01:00:00Z"), "evt-1#0")); + + TdengineQueryStatement statement = queries.telemetryFields(query); + + assertThat(statement.sql()) + .startsWith("SELECT ts, fact_id, frame_id, received_at, field_key") + .contains(" FROM tf_gb32960_") + .contains(" WHERE ts >= ? AND ts < ?") + .contains(" AND (ts > ? OR (ts = ? AND fact_id > ?))") + .contains(" ORDER BY ts ASC, fact_id ASC LIMIT ?"); + assertThat(statement.values()) + .containsExactly( + query.from(), query.to(), + query.cursor().ts(), query.cursor().ts(), query.cursor().tieBreaker(), + 100); + } } diff --git a/modules/sinks/tdengine-history-store/src/test/java/com/lingniu/ingest/tdenginehistory/TdengineJdbcHistoryReaderTest.java b/modules/sinks/tdengine-history-store/src/test/java/com/lingniu/ingest/tdenginehistory/TdengineJdbcHistoryReaderTest.java index 6aff69cf..f2b4e91d 100644 --- a/modules/sinks/tdengine-history-store/src/test/java/com/lingniu/ingest/tdenginehistory/TdengineJdbcHistoryReaderTest.java +++ b/modules/sinks/tdengine-history-store/src/test/java/com/lingniu/ingest/tdenginehistory/TdengineJdbcHistoryReaderTest.java @@ -68,6 +68,33 @@ class TdengineJdbcHistoryReaderTest { assertThat(page.items().getFirst().longitude()).isEqualTo(113.12); } + @Test + void readsTelemetryFieldRowsWithNextCursor() throws Exception { + RecordingQueryJdbc jdbc = new RecordingQueryJdbc(); + jdbc.rows.add(telemetryFieldResult("evt-1#0", "2026-06-29T05:00:01Z")); + jdbc.rows.add(telemetryFieldResult("evt-2#0", "2026-06-29T05:00:02Z")); + TdengineJdbcHistoryReader reader = new TdengineJdbcHistoryReader( + jdbc.dataSource(), new TdengineHistorySchema("vehicle_history")); + + TdenginePage page = reader.queryTelemetryFields(new TdengineTelemetryFieldQuery( + "GB32960", + "VIN123", + "VEHICLE.totalMileageKm", + Instant.parse("2026-06-29T00:00:00Z"), + Instant.parse("2026-06-30T00:00:00Z"), + TdengineQueryOrder.ASC, + 1, + null)); + + assertThat(page.items()) + .extracting(TdengineTelemetryFieldRow::factId) + .containsExactly("evt-1#0"); + assertThat(page.items().getFirst().fieldKey()).isEqualTo("VEHICLE.totalMileageKm"); + assertThat(page.items().getFirst().valueDouble()).isEqualTo(123.4); + assertThat(page.nextCursor()) + .contains(new TdenginePageCursor(Instant.parse("2026-06-29T05:00:01Z"), "evt-1#0")); + } + private static Map rawFrameResult(String frameId, String ts) { Map row = new LinkedHashMap<>(); row.put("ts", Timestamp.from(Instant.parse(ts))); @@ -113,6 +140,29 @@ class TdengineJdbcHistoryReaderTest { return row; } + private static Map telemetryFieldResult(String factId, String ts) { + Map row = new LinkedHashMap<>(); + row.put("ts", Timestamp.from(Instant.parse(ts))); + row.put("fact_id", factId); + row.put("frame_id", "frame-1"); + row.put("received_at", Timestamp.from(Instant.parse(ts).plusMillis(500))); + row.put("field_key", "VEHICLE.totalMileageKm"); + row.put("value_type", "DOUBLE"); + row.put("value_text", "123.4"); + row.put("value_double", 123.4); + row.put("value_long", null); + row.put("unit", "km"); + row.put("quality", "GOOD"); + row.put("source_path", "VEHICLE.totalMileageKm"); + row.put("raw_uri", "archive://gb32960/frame-1.bin"); + row.put("metadata_json", "{}"); + row.put("protocol", "GB32960"); + row.put("vehicle_key", "VIN123"); + row.put("vin", "VIN123"); + row.put("phone", ""); + return row; + } + private static final class RecordingQueryJdbc { private final List> rows = new ArrayList<>(); private final Map boundValues = new HashMap<>();