ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Flink + Lindorm TSDB:物联网时序数据实时存储与查询实践

Flink + Lindorm TSDB:物联网时序数据实时存储与查询实践 1. 为什么时序数据实时场景最终会走到“Flink Lindorm TSDB”这个组合我最早被时序数据“上一课”是在一个物联网监控平台上。平台接入的设备从几千台涨到几万台每个设备每 5 秒上报一次状态一天轻松攒下几十亿行数据。最初的方案简单粗暴——所有上报都写 MySQL前期确实省事设备量一上来就出问题了单表膨胀到上亿行之后写入频繁锁表查询某个设备最近 7 天的趋势曲线要花十几秒半夜的运维告警根本压不住。后来团队做过一轮完整的技术选型对比把主流的时序方案摸了一遍。开源自建的 InfluxDB 集群版成本高运维负担重TDengine 在个别版本上跟 Flink 的对接还需要打磨插件适配经常要自己写传统“ Kafka 离线数仓 ”的日级链路又满足不了实时告警的 SLA。最终在云上把方案收敛成了“Flink 流计算 Lindorm TSDB 时序存储”的组合。这个组合的逻辑并不复杂Flink 负责实时处理流动中的数据Lindorm TSDB 负责把海量时序数据存下来、还能快速查出去。前者解决“怎么算”后者解决“怎么存、怎么查”。1.1 时序数据为什么难处理时序数据有三个非常鲜明的特征直接决定了存储和计算架构的选型方向。第一个特征是只追加、不修改。传感器采集到的值一产生就是历史不需要 update只需要 insert而且写入频率极高。普通关系型数据库的索引维护、事务日志在这种写入模型下反而是负担每次插入都要维护二级索引、写 undo log这些开销在一秒几万次的写入流量下会被无限放大。第二个特征是数据随时间快速贬值。5 秒前的数据可能是实时告警的关键依据5 天前的数据只适合做趋势分析5 个月前的数据大多数场景下可以压缩到千分之一甚至直接删掉。这种“热度随时间衰减”的特性要求存储系统天生就支持按时间分区、冷热分层、高压缩比而不是靠 DBA 每隔几个月写一次归档脚本。第三个特征是查询模式高度固定。大部分查询都是“给定一个设备/一组标签查某个时间段内的若干指标”本质上就是按时间范围扫描。传统数据库的 B 树在随机点查上很强但这种大规模顺序扫描场景并不是它的主场。时序数据库普遍采用 LSM 树和列式存储写入吞吐和范围扫描性能都更有优势。1.2 Flink 在这条链路里的不可替代性很多人会问Lindorm TSDB 自己就能接收写入为什么中间非要加一个 Flink真实场景里设备上报的数据往往不能直接用。以我做过的光伏逆变器项目为例上报消息里除了标准的三相电压、电流数据还有各个厂家私有的字段设备型号、固件版本、告警码、状态位。如果不对这些消息做清洗、补全比如消息里缺了 region 维度就要根据设备注册表补上、格式统一直接入库之后查询端写 SQL 会非常痛苦。更重要的是流式数据天然存在乱序和重复的问题。传感器网络抖动导致数据延迟到达采集网关偶发重推导致同一时间点的数据推了两遍。这些都需要流处理框架做去重、乱序修正watermark 机制、滚动窗口聚合。Flink 在这些方面是当前成熟度最高的引擎之一它的事件时间处理能力和 exactly-once 状态一致性保证是自研流处理代码很难复刻的。所以这个组合的本质可以概括为Flink 把“不可控的海量乱序流”变成“质量可控的数据流”Lindorm TSDB 再把“数据流”变成“可高效查询的资产”。下面我从架构设计开始逐步展开整个集成的落地细节。2. 整条数据链路的设计与 Lindorm TSDB 模型映射2.1 从设备端到落库的数据流主线集成后的整体数据流一般长这样设备端通过 MQTT/HTTP 将采集数据上报到网关或云 IoT 平台。网关将消息转换为统一 JSON 格式写入 Kafka 作为消息缓冲。这一层不是多余的设备端的连接波动、突发流量要靠 Kafka 缓冲Kafka 也方便后续多套消费链路复用同一份数据——比如一套走实时告警另一套走持久化存储。Flink 消费 Kafka topic在流上完成 schema 解析、去重、字段补全必要时做窗口聚合比如 5 分钟平均值。处理后的数据通过 Sink 写入 Lindorm TSDB。上层应用通过 Lindorm TSDB 的 SQL 接口查询时序数据用于仪表盘、告警、算法分析。这里补充一个架构上的建议如果对实时性要求没那么苛刻可以考虑在 Kafka 和 Lindorm 之间不加 Flink直接用 Kafka Connect 或者 DataWorks 的数据同步任务把数据搬进 Lindorm。但一旦业务里有聚合计算、异常检测、多流 join 的需求Flink 就绕不开了。我的经验是先按“必然要上 Flink”的架构设计后续加计算逻辑才不会推倒重来。2.2 Lindorm TSDB 的数据模型回顾在写代码之前务必先把 Lindorm TSDB 的数据模型搞清楚。它采用“宽表 时序”模型表结构里的列分成三类tag 列、timestamp 列、field 列。tag 列用于标识数据来源和过滤条件的元数据比如 device_id、region、type。查询时的过滤条件一般都写在 tag 上Lindorm 会基于 tag 建立索引。timestamp 列数据产生的物理时间。Lindorm TSDB 按时间对数据进行分区存储时间戳必须显式声明否则无法使用时序查询能力。field 列实际监控值可以是数值型、字符串型甚至二进制。数值型字段是时序聚合查询均值、最大值、采样的主力。举个例子一张光伏设备电流表可以这样建CREATE TABLE inverter_current ( device_id VARCHAR TAG, region VARCHAR TAG, ts TIMESTAMP, current_a DOUBLE, current_b DOUBLE, current_c DOUBLE, temperature DOUBLE, PRIMARY KEY (device_id, region, ts) );注意 PRIMARY KEY 是 tag 列 timestamp 列的组合。这是一个关键设计查询时按 device_id 时间范围扫描Lindorm 能够做分区裁剪和索引查找查询性能会好很多。如果主键里没有时间字段时序数据的查询优势就发挥不出来。2.3 表结构映射中最容易埋雷的细节我在这条链路上帮人排查过不少问题见过最多的是这三个把平台接收时间当成业务时间戳。设备消息里通常有“采集时间”和“平台接收时间”两个字段有些人图省事直接用接收时间作为 ts。如果消息在 Kafka 里积压了半小时落库的时间就整体偏移趋势曲线会出现断层或错位。正确做法是使用事件时间设备采集时间Flink 侧通过 watermark 处理乱序后再落库。tag 列数量失控。为了查询方便把 firmware_version、app_version 这类低基数字段也塞进 tag。tag 会参与索引和分区数量越多写入放大越严重。建议只把查询频率高、基数可控的维度放 tag其余放 field。一张表混装多种指标。有的团队把电流、电压、温度、功率全塞进同一张表列数膨胀到几十列。更合理的做法是按指标域拆表功率类、温度类、状态类每张表保持列少而稳定方便扩展和查询优化。3. 核心实现Flink 接入 Lindorm TSDB 的三条路径3.1 先梳理清楚可用的接入方案Flink 和 Lindorm TSDB 对接目前主流有三条路径各有适用场景。先看图再决定走哪条。接入方式开发成本吞吐能力适用场景Flink JDBC ConnectorSQL低中等数据量不大、快速验证链路Lindorm 官方/托管提供的 Sink Connector中等高生产环境高吞吐写入自定义 Flink Sink基于批量写 API高高有定制逻辑、需要强控写入行为先说 Flink JDBC Connector。Flink 自带的 JDBC 连接器支持 sink 到任何支持 JDBC 的数据库Lindorm TSDB 也提供兼容 MySQL 协议的驱动。如果你只是验证链路、数据量到不了每秒几万条用这个方式最快代码量也最少。在 Flink SQL 里注册目标表CREATE TABLE lindorm_sink ( device_id STRING, region STRING, ts TIMESTAMP(3), current_a DOUBLE, current_b DOUBLE, current_c DOUBLE, temperature DOUBLE, PRIMARY KEY (device_id, region, ts) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:mysql://lindorm-xxxx:3306/iot_db, table-name inverter_current, username your_user, password your_password );然后一句 INSERT INTO 就能把处理后的数据写进去。但我必须提醒官方 JDBC 连接器在 sink 模式下的写入是批量加事务的默认按批次累积数据再提交。这在严格的实时场景下会带来两个副作用一是 checkpoint 间隔内的数据延迟可见二是数据库连接数压力大并行度一高容易把连接池打满。所以 JDBC 这条路只适合小流量验证和联调。3.2 自定义 Sink 实现高吞吐写入生产场景我更推荐自定义 Sink思路是在 Flink 的 RichSinkFunction 里维护一个内存 buffer攒够一批或者到达时间窗口后通过 Lindorm 的批量写接口一次性提交。这样做的好处批量写减少了网络 RTTLindorm 的批量接口本来就是为了百万级写入设计的同时我们能精确控制 flush 时机避免 checkpoint 期间的连接风暴。下面给一个简化但能跑通的示例逻辑public class LindormTsdbSink extends RichSinkFunctionMonitorRecord { private transient LindormTSDBClient client; private transient ListMonitorRecord buffer; private static final int BATCH_SIZE 1000; private static final long FLUSH_INTERVAL_MS 2000; Override public void open(Configuration parameters) throws Exception { client LindormTSDBClientFactory.create(endpoint, username, password); buffer new ArrayList(BATCH_SIZE); } Override public void invoke(MonitorRecord value, Context context) throws Exception { buffer.add(value); if (buffer.size() BATCH_SIZE) { flush(); } } private void flush() throws Exception { if (buffer.isEmpty()) return; client.batchWrite(buffer); buffer.clear(); } Override public void close() throws Exception { flush(); if (client ! null) client.close(); } }示例里省掉了很多生产必需的细节实际落地时一定要补上失败重试与死信处理batchWrite 失败后要带退避重试重试仍然失败的批次不能直接丢弃要进入死信队列比如写回 Kafka 的 dead topic方便事后排查。与 checkpoint 联动如果要求端到端一致性则需要实现两阶段提交或者在写入前做幂等去重这块下一章展开讲。定时 flush示例里只做了按条数触发但源数据流量低时可能长时间不触发数据延迟变大。应该加一个 timer到点不管够不够条数都 flush。3.3 关于官方连接器的使用建议如果你的环境里能用 Lindorm 官方或云厂商托管提供的 Flink Sink 连接器那当然最好——毕竟少造轮子。即便直接用官方连接器也要注意几个配置点直接影响写入效果sink.buffer-flush.max-rows触发批量写入的条数。太大会让数据在内存里停留过久太小又达不到批量效果。我实践下来通常设在 1000~5000 之间。sink.buffer-flush.interval批量写入的最大等待时间。建议 1~3 秒比 checkpoint 间隔短一些这样 checkpoint 成功后数据基本已经落库不容易出现“checkpoint 成功但外部存储迟迟没有数据”的伪成功现象。sink.max-retries单批次写失败后的重试次数。重试必须配合退避策略否则一次故障会瞬间把下游打爆。4. 生产环境里的写入调优、一致性保障与成本控制4.1 调优从理解三个瓶颈开始Flink 写 Lindorm TSDB瓶颈通常不在 Flink 引擎本身而在“内存 buffer 积压”“下游写入吞吐”和“网络连接数”这三处。我逐个说一下定位方法。内存 buffer 积压的典型表现是 Flink Web UI 的背压指标打红。每次 flush 后没有及时释放 buffer或者批量条数设置过大比如 1 万条导致单次 flush 耗时很长sink 算子就卡住了。解法是调小 batch size、缩短 flush interval同时给 sink 算子单独调大并行度。下游写入吞吐的问题可以在 Lindorm 控制台观察写入 QPS 和写入延迟。如果 Lindorm 本身出现限流或者延迟升高大概率是集群规格不够或者表分区设计不合理。时序表的分区键建议按高频查询维度设计让写入流量均匀分布避免出现热分区。网络连接数方面Flink 每个并行子任务都会维护到 Lindorm 的长连接。并行度从 4 调到 32连接数也随之上涨。如果连接池参数没调超过上限会导致创建连接超时。我的经验是让单个连接的吞吐高一些而不是开一堆空闲连接占资源。4.2 端到端一致性幂等性才是关键很多人在 Flink 集成外部存储时会纠结 exactly-once 能不能实现。Lindorm TSDB 的写入本质上是 upsert 语义——同一个主键tag timestamp重复写入是幂等的后写的覆盖先写的。这恰好让 Flink checkpoint 重放机制变得可行哪怕任务重启导致同一批数据被发送两次只要主键一致落库结果不会翻倍。这是时序数据库做实时链路的一个天然优势。但这里有个前提写入必须按主键幂等。如果你在 Flink 里做了窗口聚合输出的是带窗口起止时间的新记录你需要把窗口时间作为主键的一部分而不是用当前物理时间否则窗口重算后写的是不同时刻的数据就会产生重复记录。4.3 压缩、精度与生命周期成本控制的核心Lindorm TSDB 的存储成本大致由三个维度决定数据精度、保留周期、压缩策略。分享一下我们调参的实际经验精度float 尽量用 4 字节单精度不要因为源数据里带了 6 位小数就全部上 double。时序数据量大每一列宽一点总量差就是成倍的。保留周期按业务热度拆成多条生命周期策略。比如热数据保留 7 天全精度7~30 天降采样成 1 分钟均值30 天以上的数据按 5 分钟聚合存储。Lindorm 支持冷热分层和 TTL这类策略可以自动生效运维省很多事。压缩率时序值的重复度高恒定电压、缓慢变化的温度Lindorm 的列存压缩对这些模式效果很好。但如果你在表里塞入大量低基数字符串 tag压缩率会被严重稀释——这也是我前面强调 tag 要收敛的原因之一。5. 集成实战中的典型问题与排查链路5.1 时间戳类型不匹配最常见的翻车现场用 Flink SQL 往 Lindorm 写数据时最容易翻车的是时间戳精度不一致。Flink 的 TIMESTAMP(3) 默认精确到毫秒但很多设备端的采集时间是秒级。如果你直接用事件时间字段作为 ts会出现 1000 倍的时间偏移查询出来的曲线整个变形。分享一个排查技巧先不急着推全量在小流量模式下写入一批数据然后到 Lindorm 侧执行一条查询对比落库时间戳和当前系统时间。如果差了 1000 倍十有八九是秒和毫秒没对齐。处理方式是在 Flink SQL 里显式转换SELECT device_id, region, TO_TIMESTAMP_LTZ(event_time, 0) AS ts, -- 将秒级时间转成毫秒级 current_a, current_b, current_c FROM source_table5.2 schema 变动导致反序列化失败IoT 场景里厂商升级固件之后上报的 JSON 多塞几个字段是家常便饭。Flink 消费者如果用了强类型 schema新字段没有提前注册会导致反序列化失败整个作业卡住。这个问题有三层解法。第一层消息进入 Kafka 之前先做一层格式规整把关键字段补齐第二层Flink 消费 JSON 时使用 ignore-parse-errors 选项保证脏数据不阻断主链路第三层解析失败的数据通过侧输出流单独写入一个死信 topic留作后续排查。我们线上就是靠第三层挽救了很多排查线索——一旦业务方反馈数据不对先查死信 topic往往能直接定位问题源头。5.3 背压问题的定位路径任务跑着跑着发现 Lindorm 写入 QPS 并不高但 Flink 作业背压持续红色。这种情况多数不是 Lindorm 的问题而是 sink 算子内部 flush 慢。我的定位路径是这样走的看 Flink Web UI 的背压面板确认是哪个算子背压通常集中在 sink 算子。看 sink 算子的 busy 时间占比如果接近 100%说明算子一直在工作而不是等待。去 Lindorm 控制台看该时段的写入请求延迟和拒绝数判断是否有服务端限流。对照日志时间看是不是某次 batch 返回超时导致大量重试积压。绝大多数情况下把批量大小调小、重试退避参数优化一下背压就消退了。5.4 “丢数据”的真相往往是“晚到数据”经常有业务方跑来说数据丢了查下来不是丢而是数据乱序导致被窗口丢弃。Flink 默认在 watermark 之后到达的数据会被视为迟到数据如果窗口已经触发计算迟到数据默认被丢弃。处理这个问题有两个关键点根据设备实际生产环境设置合理的allowedLateness参数。比如光伏项目设备数据延迟最大不超过 2 分钟就设置成 2 分钟。开启 side output把迟到的数据单独输出由下游决定是补算还是忽略。而不是默默丢弃否则数据质量监控完全黑屏。我在自己的项目里始终在 Flink 作业里保留一条迟到数据计数监控。如果某天迟到数据量异常飙升大概率是设备端网络或者采集网关出了问题这时候提前预警比事后查数据要省力得多。5.5 连接数打满导致作业雪崩还有一个高并发下容易踩的坑Flink 作业并行度调高之后每个并行子任务都建立独立连接到 Lindorm加上多个作业共享同一个 Lindorm 实例连接数很容易打满。表现是作业启动正常跑一段时间后开始频繁报连接超时重启恢复过一会又超时。解决思路有两个方向一是收敛并行度让单个连接的吞吐上来而不是靠连接的绝对数量堆吞吐二是在 Flink 侧做连接复用比如在 RichSinkFunction 的 open 方法里只初始化一个 clientsink 算子内部用线程安全的方式共享。不要指望连接池的默认参数能适配所有场景生产环境务必压测一轮连接数和吞吐的关系。6. 写在最后的实操心得把 Flink 和 Lindorm TSDB 打通这件事难度其实不在“连上”而在“连好”。连上一个 JDBC 接口可能只要半小时但真要支撑千万级设备、几十亿日增数据的实时链路模型设计、参数调优、异常链路处理每一环都藏着一堆经验活。我个人的体会有三点第一表模型设计一定要在写第一行代码之前就定清楚tag 收敛、按指标域拆表、时间精度统一这三件事做对了后面少走一半弯路第二流的质量远远比计算逻辑重要脏数据、迟到数据、重复数据这三种情况要在一开始就设计好处理策略而不是等出问题再救火第三上线前务必压测特别是连接数、批量大小和背压的关系这些在文档里永远查不到只能靠模拟流量试出来。如果你正在规划一条新的时序数据实时链路我建议先拿小流量把全链路跑通再从查询侧反推表结构设计最后才是调参。这个顺序能帮你避开很多“做完才发现查不动”的悲剧。
返回列表