档案放 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 也对应建好:metadataJdbcTemplate和tdengineJdbcTemplate。前者查档案,后者查时序。注入时用@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,一个查询里拿不到两张表。怎么办?服务层编排。比如查询设备最新状态:
- 先从 PG 查设备档案
- 再从 TDengine 查最新遥测点
- 服务层把两者拼装成
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 层拼装。
回顾一下,双数据源查询层的核心就这么几件事:
- 两个 DataSource + 两个 JdbcTemplate,@Primary + @Qualifier 区分
- 白名单防注入:能参数绑定的全绑定,不能绑定的锁白名单
- limit+1 分页,多查一条判断 hasMore
- 在线判定:最新点 + 5 分钟阈值;离线用 HAVING 聚合后过滤
- 健康检查:TDengine 挂了 API 照样跑
你在实际项目里遇到过双数据源的坑吗?比如连接池耗尽、事务失效、或者 SQL 注入的隐患?欢迎在评论区聊聊你的方案。
觉得有用?点个关注,持续获取优质内容。