feat: query telemetry fields from tdengine

This commit is contained in:
lingniu
2026-06-29 13:55:25 +08:00
parent 84de9b2dfa
commit 90342fe5bf
7 changed files with 166 additions and 0 deletions

View File

@@ -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,

View File

@@ -7,4 +7,6 @@ public interface TdengineHistoryReader {
TdenginePage<TdengineRawFrameRow> queryRawFrames(TdengineRawFrameQuery query) throws IOException;
TdenginePage<TdengineLocationRow> queryLocations(TdengineLocationQuery query) throws IOException;
TdenginePage<TdengineTelemetryFieldRow> queryTelemetryFields(TdengineTelemetryFieldQuery query) throws IOException;
}

View File

@@ -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<TdengineTelemetryFieldRow> 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 <T> TdenginePage<T> readPage(TdengineQueryStatement statement,
int pageSize,
SqlRowMapper<T> 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> {
T map(ResultSet rs) throws SQLException;

View File

@@ -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);
}
}

View File

@@ -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);
}
}

View File

@@ -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<TdengineTelemetryFieldRow> 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<String, Object> rawFrameResult(String frameId, String ts) {
Map<String, Object> row = new LinkedHashMap<>();
row.put("ts", Timestamp.from(Instant.parse(ts)));
@@ -113,6 +140,29 @@ class TdengineJdbcHistoryReaderTest {
return row;
}
private static Map<String, Object> telemetryFieldResult(String factId, String ts) {
Map<String, Object> 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<Map<String, Object>> rows = new ArrayList<>();
private final Map<Integer, Object> boundValues = new HashMap<>();