【系列:TDengine 工业物联网实战:从零搭起可运行系统 · 第 8 篇】
2026/8/20 23:57:48 网站建设 项目流程

档案放 PostgreSQL,时序放 TDengine——这是很多 IoT 系统的标准姿势。但真要把两个数据源塞进同一个 Spring Boot 应用,配置、注入、防注入、分页、健康检查,每一步都有坑。本文用实战源码,讲透双数据源查询层的完整实现。

数据写进去了,怎么查出来?

前几篇我们把 MQTT 消息落进了 TDengine 超级表,设备档案存进了 PostgreSQL。现在问题来了:一个 Spring Boot 应用,怎么同时连两个数据库?

你可能会想:搞个分布式事务?引入中间件?都不需要。查询层的做法,是让两个数据源各司其职,服务层编排。今天这篇,我们就从零拆解java-api/模块的双数据源实现。

双数据源配置:4 个 Bean 打天下

先看核心配置类DataSourceConfiguration.java。这里定义了 4 个 Bean:两个 DataSource,两个 JdbcTemplate。

**档案数据源(PG)**用 Hikari 连接池,标了@Primary

@Bean("metadataDataSource")@Primary@ConfigurationProperties("spring.datasource.hikari")publicHikariDataSourcemetadataDataSource(@Qualifier("metadataDataSourceProperties")DataSourcePropertiesproperties){returnproperties.initializeDataSourceBuilder().type(HikariDataSource.class).build();}

**时序数据源(TDengine)**用的是DriverManagerDataSource,每次新建连接,不给连接池:

@Bean("tdengineDataSource")publicDataSourcetdengineDataSource(TdenginePropertiesproperties){DriverManagerDataSourcedataSource=newDriverManagerDataSource();dataSource.setDriverClassName("com.taosdata.jdbc.ws.WebSocketDriver");dataSource.setUrl(properties.url());dataSource.setUsername(properties.username());dataSource.setPassword(properties.password());returndataSource;}

注意这个驱动:com.taosdata.jdbc.ws.WebSocketDriver。它走的是TAOS-WS 协议,连的是 6041 端口——正是第 6 篇里那个 WebSocket 门。这意味着 Java 侧不需要装本地客户端,一个 JDBC 驱动就搞定。

URL 长这样:

jdbc:TAOS-WS://localhost:6041/iot?user=root&password=taosdata

两个 JdbcTemplate 也对应建好:metadataJdbcTemplatetdengineJdbcTemplate。前者查档案,后者查时序。注入时用@Qualifier区分——比如 Repository 里是构造器注入:

publicTelemetryRepository(@Qualifier("tdengineJdbcTemplate")JdbcTemplatejdbc){this.jdbc=jdbc;}

TDengine 侧没配连接池DriverManagerDataSource每次查询新建连接。查询型接口的负载可控时这样最省事,连接数压上去再换 Hikari 也不迟。

时序模板额外设了两件事:setQueryTimeout(30)(慢查询 30 秒掐断)和setFetchSize(1000)(流式取数,防止大结果集撑爆内存)。

为什么档案和时序必须分开?

有人问:都放一个库里不行吗?

看数据特征就明白了。

PG 里的 device 表:id、display_name、product_key、factory_id、workshop_id、device_type、region、enabled、metadata JSONB。这是档案数据——小、低频变更、事务性强。改名、停用、改 metadata,都是典型的关系库操作。

TDengine 里的 telemetry 超级表:温度、湿度、电压、电流……这是时序数据——大、只追加、按时间窗口聚合。用 TDengine 的超级表 + 窗口函数,一条 SQL 就能算小时均值。

关键约束:不跨库 join。档案在 PG,时序在 TDengine,一个查询里拿不到两张表。怎么办?服务层编排。比如查询设备最新状态:

  1. 先从 PG 查设备档案
  2. 再从 TDengine 查最新遥测点
  3. 服务层把两者拼装成DeviceLatest返回

查询白名单:列名和 INTERVAL 怎么防注入?

这是本篇的一个关键点。

看这段聚合 SQL:

Stringsql=""" SELECT _wstart AS window_start, AVG(%s) AS average_value, ... FROM iot.telemetry WHERE device_id = ? AND ts >= ? AND ts < ? INTERVAL(%s) """.formatted(safeMetric,safeMetric,safeMetric,safeMetric,safeInterval);

注意:聚合列名和 INTERVAL 是拼进去的,但设备 ID 和时间范围全用?参数绑定。

为什么列名不能参数绑定?因为 SQL 语法上,列名和窗口函数参数不是值,?占位符在这里不生效。你没法写AVG(?)让 JDBC 帮你填列名。

那怎么办?白名单校验

privatestaticfinalSet<String>METRICS=Set.of("temperature","humidity","voltage","current_value","power","pressure","flow_rate","rotational_speed","vibration");

METRICS 白名单 9 项,INTERVALS 白名单 7 项:

10s / 30s / 1m / 5m / 15m / 1h / 1d

请求里的 metric 参数先进requireMetric()校验,不在白名单直接抛IllegalArgumentException。校验通过后才用.formatted()拼进 SQL。

设备 ID 和时间范围呢?全部?绑定。Timestamp.from(start)直接传参,JDBC 驱动处理转义。

这就是双保险:能参数绑定的绝不拼字符串,不能绑定的用白名单锁死

limit+1 分页:多查一条比多查一页便宜

分页查询有个经典问题:怎么知道还有没有下一页?

常规做法是再查一次 count。这里用的是limit+1

List<TelemetryPoint>rows=repository.findTelemetry(deviceId,start,end,safeOffset,safeLimit+1);returnPageResponse.of(rows,safeOffset,safeLimit);

limit + 1条,如果返回的行数大于请求的 limit,说明还有更多数据:

booleanhasMore=rows.size()>requestedLimit;List<T>items=hasMore?rows.subList(0,requestedLimit):rows;

多查一条比多查一页便宜——省掉一次 count 查询,代价只是一条记录的传输。

PageResponse四个字段:items / offset / limit / hasMore。

配套的QueryRangeValidator做参数校验:

  • start/end 必填且 start < end
  • 跨度 ≤ maxRangeDays(默认 31 天)
  • limit 1~5000
  • offset ≥ 0

违规直接抛InvalidQueryException

在线状态:5 分钟阈值 + HAVING 过滤

设备在线怎么判定?看最新一条遥测的时间戳。

privatestaticfinalDurationONLINE_THRESHOLD=Duration.ofMinutes(5);booleanonline=point!=null&&point.timestamp().isAfter(Instant.now().minus(ONLINE_THRESHOLD));

最新点时间在 now-5min 内算在线,否则离线。阈值固定 5 分钟,这是业务规则,写死在常量里。

离线设备列表呢?看这条 SQL:

SELECTdevice_idFROMiot.telemetryGROUPBYdevice_idHAVINGLAST(ts)<?LIMIT?

注意是 HAVING 不是 WHERE。为什么?WHERE 在分组前执行,这时候 LAST(ts) 还没算出来。HAVING 在分组后过滤,才能用聚合函数的计算结果。

/api/operations/offline-devices?minutes=30&limit=这个接口,minutes 参数限制 1~43200(30 天),防止有人传个负数或超大值把 TDengine 压垮。

PG 档案 CRUD:JSONB 和 RETURNING 的妙用

档案侧的几个细节值得说。

metadata 是 JSONB 类型,插入时用 ObjectMapper 序列化后 CAST:

INSERTINTOdevice(...)VALUES(?,...,CAST(?ASjsonb))RETURNING*

RETURNING 一步拿回更新后的行

UPDATEdeviceSET...WHEREid=?RETURNING*

如果返回 null,说明设备不存在,抛NotFoundException。省了一次 select。

可选过滤用CAST(? AS VARCHAR) IS NULL技巧

WHERE(CAST(?ASVARCHAR)ISNULLORfactory_id=?)

参数为 null 时条件恒真,一个 SQL 搞定可选过滤,不用动态拼 SQL。

异常处理上,DuplicateKeyException转成ConflictException,返回 409 语义。

健康检查与可运维性

TDengine 挂了,API 不能挂——这是设计原则。

TdengineHealthIndicator做探针:

SELECTSERVER_VERSION()

成功即 UP,异常 DOWN。Kubernetes 的 liveness/readiness 探针可以直接用这个端点。

其他配置:

  • server.port 默认 8080,JAVA_API_PORT环境变量可覆盖
  • shutdown: graceful,30 秒排空
  • compression 开启
  • management 暴露 health/info/prometheus/metrics
  • springdoc 开启,/swagger-ui.html 在线文档

接口全景与双库协作边界

最后过一遍 8 个 Controller 的接口清单:

方法+路径说明
POST /api/devices创建设备(metadata 任意 JSON)
GET /api/devices?factoryId=&deviceType=&offset=&limit=分页列表(可选过滤)
GET /api/devices/{id}档案详情
PUT /api/devices/{id}更新(含 enabled 停用)
GET /api/devices/{deviceId}/latest最新遥测 + online 布尔
GET /api/devices/{deviceId}/telemetry?start=&end=&offset=&limit=分页遥测
GET /api/devices/{deviceId}/aggregate?metric=&interval=&start=&end=窗口聚合
GET /api/factories/{factoryId}/statistics?metric=&start=&end=工厂级汇总
GET /api/vehicles/{vehicleId}/track?start=&end=&offset=&limit=车辆轨迹分页
GET /api/operations/offline-devices?minutes=&limit=离线设备列表

注意 alarm-rules 相关接口只列了名,那是第 9 篇的素材——告警引擎会展开讲。

双库协作的边界:不 join、不跨库事务、服务层编排。档案查 PG,时序查 TDengine,数据在 Service 层拼装。


回顾一下,双数据源查询层的核心就这么几件事:

  1. 两个 DataSource + 两个 JdbcTemplate,@Primary + @Qualifier 区分
  2. 白名单防注入:能参数绑定的全绑定,不能绑定的锁白名单
  3. limit+1 分页,多查一条判断 hasMore
  4. 在线判定:最新点 + 5 分钟阈值;离线用 HAVING 聚合后过滤
  5. 健康检查:TDengine 挂了 API 照样跑

你在实际项目里遇到过双数据源的坑吗?比如连接池耗尽、事务失效、或者 SQL 注入的隐患?欢迎在评论区聊聊你的方案。


觉得有用?点个关注,持续获取优质内容。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询