物联网时序数据同步实践:从TDengine到Doris的完整链路
开头就直接切入吧。我去年在团队里搭过一条物联网数据链路设备产生的时序数据落在 TDengine 里运营和算法那边需要做多维分析又不敢让几十张报表直接压到时序库上最后决定把全量增量数据都同步到 Doris。任务调度用的是 Apache DolphinScheduler整个链路由 AllData 数据中台统一管理。这篇文章就把这条“TDengine 到 Doris”的同步全流程完整复现一遍包括为什么这么选型、读写两端的表设计和参数怎么定、调度和告警怎么配以及我当时踩过的几个比较隐蔽的坑。如果你正在搭类似的数据同步任务或者被“时序库里的数据怎么稳定搬到 OLAP 库”这个问题卡住可以参考这条链路。1. 这个同步任务到底在解决什么问题为什么物联网数据必须走这条链路1.1 一条典型的物联网数据链路长什么样物联网场景的数据链路和传统业务系统有个很大的区别写入压力大、单条价值低、时间维度极强。设备侧通过 MQTT、Modbus、OPC UA 等协议上报数据经网关做协议解析后写入 TDengine。TDengine 负责承接高频写入和近期数据的快速查询比如“过去一小时某台设备的温度曲线”。但真正的分析需求往往更复杂——要看某个区域所有设备的综合指标、要跟业务台账关联、要做任意时间跨度的大规模聚合这时候时序库就不是最优解了。我们的目标是让数据“既有时序库的高性能写入又有 OLAP 分析引擎的大规模查询能力”。实现方法就是再做一份数据到 Doris。听起来很简单真正做起来要解决的问题很多同步的节奏怎么定是全量一次性搬运还是按时间窗口持续增量数据在传输过程中怎么保证不丢不重同步任务失败了怎么办怎么发现、怎么自动恢复调度系统怎么跟数据源、目标库的状态联动。这就是本文要拆解的核心问题。1.2 为什么分析侧不直接查 TDengine非要再搬一份数据到 Doris先说结论TDengine 和 Doris 解决的不是同一类问题。TDengine 是时序数据库强在写入吞吐、时间范围扫描、按时间窗口的聚合INTERVAL以及超级表 子表的标签查询模式。但它在复杂多表关联、大并发即席查询、高基数去重统计这些方面的表现和 Doris 这种 MPP 列式分析引擎有明显差距。Doris 存的是列式数据向量化执行引擎应付大数据量聚合非常稳而且支持 MySQL 协议BI 工具接起来很容易。数据冗余存两份不是浪费而是隔离。把分析查询放到 Doris 上和生产写入到 TDengine 之间就解耦了——报表再慢再重也影响不到设备数据采集这条主链路。这也是我后来跟团队反复强调的一个点同步任务不是简单的“复制数据”而是给分析场景建立一个独立的数据消费渠道。理解了这一点后面很多架构取舍就自然了。1.3 这个同步任务的要求清单做之前先列需求不然做到一半容易跑偏。我们的核心要求是这么几条增量为主、全量兜底日常每 5 分钟拉取一次 TDengine 中新增的数据数据模型变更或口径调整时支持手动触发全量重刷。不重不丢同步任务必须有明确的位点概念记录每次拉取到的最大时间戳下次从位点之后继续读同时导入 Doris 时要按批次做标签label去重防止网络重试造成重复写入。可视化编排与失败重试通过 DolphinScheduler 的 DAG 工作流维护整个流程失败自动重试重试仍失败则触发告警。告警必须能触达值班人员企业微信机器人告警是我们当时的要求任务失败 1 分钟内就要收到通知。这些要求听起来都是常识但越常识的东西越容易在实现时被忽略。后面我会逐个讲每个要求是怎么落到具体配置上的。2. 技术选型与整体链路为什么是 TDengine Doris DolphinScheduler2.1 调度框架为什么选 DolphinScheduler它比别的方案强在哪调度框架的选择决定了这个同步任务能不能长期稳定跑。我们当时对比过 Airflow、XXL-Job 和 DolphinScheduler。Airflow 生态好但部署重Python DAG 有学习成本XXL-Job 定时和失败重试做得不错但不擅长表达复杂的任务依赖关系。DolphinScheduler 最突出的点是原生支持工作流 DAG——一个同步任务可以被拆成“检查数据源连通性 → 读取数据 → 导入 Doris → 校验数据量 → 更新位点”多个节点节点间依赖关系在界面上直接拖出来一目了然。实际用下来还有几个隐形的优势DolphinScheduler 中文文档和社区活跃度高遇到问题能搜到真实的解决方案它自带 worker 分组机制可以把重的同步任务调度到指定的执行节点上避免和其他业务任务互相抢资源告警方式内置了企业微信、钉钉、邮件等渠道不需要我们自己写通知模块。这些能力叠加起来让我觉得它就是做数据同步编排的正解。2.2 数据搬运怎么做自研脚本 vs DataX vs 直接写 StreamLoad调度框架定了接下来是数据搬运方式。一句话总结TDengine 侧读取Doris 侧导入中间的选择题很多。TDengine 提供了 JDBC 驱动、REST API、各种语言的连接器Doris 最推荐的导入方式是 StreamLoad走 HTTP 协议支持 CSV、JSON 格式通过 MySQL 协议也能执行 INSERT。当时我评估了三种搬运方案自研 Python 脚本用 taos 连接器读取 TDengine拼好数据后调用 Doris 的 StreamLoad HTTP 接口。优点是灵活可以精确控制位点、批次大小和错误处理逻辑缺点是要自己处理异常、重试、监控代码量不小。DataX 离线同步DataX 的插件机制支持非常多的数据源配上定时调度就能跑。优点是配置化程度高但我们的增量位点管理比较定制化DataX 天然支持得不够细最终还是需要包一层。平台自带功能AllData 里集成了数据集成能力数据源、任务配置和调度可以统一管理底层还是通过脚本或者 DataX 去执行对我们来说减少了重复造轮子的工作量。最终我们选了“自研脚本 DolphinScheduler 调度 AllData 管理元数据”的组合。原因很简单增量同步的核心复杂度在“水位线管理”这个必须自己掌控而调度、告警、依赖这些通用能力应该交给成熟组件不用自己写。下面的章节我会把每一步的关键细节都展开。2.3 AllData 在整个链路里承担什么角色AllData 在项目里更多是作为统一入口存在的。数据源连接信息、同步任务版本、目标表结构、依赖的调度工作流都在 AllData 里统一管理。DolphinScheduler 负责底层任务的真正执行AllData 负责把“哪张表从哪来到哪去、用哪个调度、数据量是多少”这类元数据串起来。这样做最直接的好处是团队里任何一个人打开 AllData 就能看到这条同步链路全貌而不是去翻一堆散落的脚本和 cron 表达式。如果哪一天不再用 DolphinScheduler切换成本也只是改一层适配而不是推倒重来。中台产品听起来“重”但当你管理的同步任务数量超过 10 个以后这种统一管理的收益就会非常明显。3. TDengine 侧数据读取与超级表设计要点3.1 超级表和子表的设计从一张普通业务表说起TDengine 和传统关系型数据库的表模型完全不同。如果你从 MySQL 迁移一批设备数据过来最自然的做法就是为每种设备类型建一个超级表STable再用每个设备建一张子表。超级表定义统一的 schema包括普通字段和标签字段子表继承超级表结构并且有自己的标签值。举个例子CREATE STABLE meters ( ts TIMESTAMP, voltage FLOAT, current FLOAT, power FLOAT ) TAGS (device_id BINARY(32), location BINARY(64)); CREATE TABLE meter_001 USING meters TAGS (device_001, shenzhen); CREATE TABLE meter_002 USING meters TAGS (device_002, shanghai);这样做的好处非常明显想查某个具体设备直接走子表速度快想按标签聚合“所有深圳设备的平均电压”查询超级表加上标签过滤TDengine 会自动在所有子表上并行扫描。这套模型和“每台设备上报数据”的物联网场景几乎是天生的匹配。如果你现在手里还只有一张带 device_id 字段的 MySQL 表要转成超级表思路就是除了主键时间戳和测点字段外的描述性字段几乎都应该挪到 TAGS 里。3.2 时序一致性为什么多表查询不会乱序热词里有人问“tdengine 如何做到多个表时序一致”这个问题其实是在问时序数据的查询语义。TDengine 的每张表主键都是时间戳所有数据严格按照时间排序落盘。因为标签字段的存在不同设备同一时间点的数据各自存储在独立的子表里但它们的物理时间戳是同一个坐标系下的——天然可比。真正要注意的不是存储而是聚合查询里的对齐语义。比如“计算所有设备最近 24 小时每 5 分钟的电压最大值”要写INTERVAL(5m)TDengine 会把所有子表的时间戳对齐到相同的窗口产出每个窗口的聚合结果。窗口对齐规则是固定的按起始时间对齐不是每张表各对各的这样下游拿到 Doris 里的数据才是时序一致的。我们当时还专门在同步任务里加了一步校验对比源库和 Doris 端按同样时间窗口计算出来的记录数不一致就告警。3.3 增量读取的正确姿势水位线、游标和批量拉取增量同步最核心的就是水位线watermark。我们定的水位线是 TDengine 里每张超级表的最大时间戳MAX(ts)。同步任务每次执行时做这几件事从 Doris 的状态表里读出上次同步的最大时间戳 last_max_ts。查询 TDengineSELECT ... WHERE ts last_max_ts AND ts 当前时间 - 延迟容忍窗口。将查询结果批量写入 Doris同一批用固定 label 标记。成功后把本次查询到的最大时间戳更新到 Doris 状态表。延迟容忍窗口是什么设备数据写入 TDengine 后偶尔会因为网络或网关缓冲出现少量延迟到达的数据。如果同步任务永远只读取“当前时间之前已经稳定落库”的数据就可以避免把迟到数据漏掉。我们用的是 5 分钟窗口代价是 Doris 里的数据比 TDengine 最多晚 5 分钟但对分析场景完全够用。另一个注意点是分页查询。TDengine 的数据量是按亿计的千万不要用LIMIT offset, size深翻页效率极差而且非常容易查询超时。正确做法是按时间维度分段拉取每次查询限制在 1 小时到 6 小时的时间范围内通过主键时间戳本身就天然有序段内即使数据量很大也可以顺序消费完。我们的脚本每次循环取 10 万行边取边导内存占用非常稳定。3.4 写入后立刻读取临时数据会不会读不到很多人问过“tdengine 保存临时数据马上读取”的问题。实际上 TDengine 的写入在客户端收到成功的确认响应后数据对读路径就是可见的不需要像某些大数据组件那样等待刷新。这也意味着同步任务在读取时并不需要担心“刚写入还没 flush 导致读不到”。真正要注意的是不要在写入峰值时刻做大规模查询因为读写都在争抢资源可能互相拖累。把这个问题放到同步任务设计里就是调度时间要避开写入高峰。我们这边设备上报高峰集中在整点前后 10 分钟所以同步任务都错峰到整点后 15 分钟启动。这样既保证了读到的数据足够全又不会给 TDengine 的写入造成明显压力。4. Doris 侧表模型、分桶策略与导入前准备4.1 集群部署最基本的结构FE 和 BE 各司其职Doris 的架构很多人听过但一上手就晕其实是两个角色的分工没搞清楚。FEFrontend负责元数据管理、SQL 解析和查询计划生成BEBackend负责数据存储和查询执行。部署上最少是 1 FE 1 BE 就能跑但生产环境一般建议 3 FE 3 BE 起步FE 之间通过选举保持元数据一致BE 之间通过副本机制保证数据冗余。部署要注意几点FE 和 BE 之间网络要通千万不能用公网 IP 跨机房部署延迟一高整个集群的性能预期就没法保证了FE 的元数据目录和 BE 的数据目录都必须放到独立磁盘不能和系统盘混在一起Doris 自身会做副本修复所以不需要在应用层做复杂的双写。如果你只是测试环境或者数据量只有几十 GB1 FE 1 BE 也能用但分区副本数设成 1 就够别默认配 3否则会一直报副本不健康。4.2 针对时序数据应该选哪种表模型Doris 有三种表模型Duplicate、Unique、Aggregate理解它们对建表特别重要。Duplicate 模型保留明细数据重复行可以同时存在适合存储原始事件或日志Unique 模型会按指定 key 列去重保留最新一条记录适合做更新类数据Aggregate 模型会在数据导入时按聚合函数预先汇总适合结果已经明确是大粒度汇总的场景。我们同步 TDengine 的时序数据选的是 Duplicate 模型。原因很直接时序数据本身就是不可变的 append-only 数据流不存在“按业务主键更新”的需求如果同一时间点重复导入了两次Duplicate 模型下我们会通过导入批次 label 和下游校验去发现重复而不是靠数据库去重静默吞掉。这样反而更利于暴露同步逻辑的问题。时间列ts和测点字段都建成普通列再在 ts 上建按天分区查询时能直接裁剪分区效率高很多。4.3 只有几 MB 数据也要分桶吗这可能是这次分享里最容易被忽略的一个问题。很多人觉得“数据才几 MB分桶完全多余一个桶就好了”听起来没错但放在 Doris 里不一样。分桶BUCKETS主要不是按数据量设计的而是按“查询并发”和“数据分布”设计的。即使表里只有几 MB 数据如果我们完全不分桶那么整张表只有一个 tablet每个查询只能由一个 BE 上的一个线程处理完全没办法利用多 BE 并行能力。将来数据增长想再改分桶是要重新建表的迁移成本很高。我们的策略是默认建表时按设备 ID 的 hash 分桶桶数按“未来 3 个月的数据增量”估算而不是按当前数据量。比如预估日增 1 亿行、单桶承载 2000 万行左右50 个桶就是合理选择。而对于几个 MB 的小维表不是不分桶而是分少量桶比如 3 个并且开副本重点在“保障查询稳定性”而不是桶多。分桶字段的选择也要注意。时序数据同步到 Doris 后最常见的查询是按设备维度过滤所以桶字段选 device_id 这种标签字段最合适如果选了时间戳做分桶字段时间范围查询会退化成所有桶全扫效果很差。4.4 建表时把字段注释写清楚后面能少吵很多架这个算是我们在实际协作里吃了亏才补上的规则。同步任务涉及三方团队设备组维护 TDengine 的超级表数据组做同步和调度分析组消费 Doris 的数据写报表。如果没有字段注释消费方对“power 到底是瞬时功率还是日累计电量”“voltage 的单位是 V 还是 mV”这类问题全靠猜。后来我们要求 Doris 建表时每个字段都必须写 COMMENT上游 TDengine 侧也要求超级表带注释字段并且在同步管道的元数据里做字段映射的说明。CREATE TABLE ods_device_meters ( ts DATETIME COMMENT 采集时间精度毫秒, device_id VARCHAR(32) COMMENT 设备编号对应TDengine子表标签, voltage DOUBLE COMMENT 电压单位V, current DOUBLE COMMENT 电流单位A, power DOUBLE COMMENT 瞬时功率单位W ) COMMENT 设备测点明细数据来自TDengine超级表meters DUPLICATE KEY(ts, device_id) PARTITION BY RANGE(ts) (...) DISTRIBUTED BY HASH(device_id) BUCKETS 50 PROPERTIES (replication_num 2);别看这只是一行行注释等任务跑起来以后分析团队能自己看懂表结构不会每天拿相同的问题轰炸数据团队。这种事属于不做没人提、做了大家都省心的典型。5. DolphinScheduler 工作流编排任务、参数、告警5.1 工作流定义与任务类型怎么选DolphinScheduler 里一个工作流可以包含多个任务节点任务之间可以配置依赖关系。我们这条同步链路的工作流拆成几个阶段先是一个前置检查节点探测 TDengine 和 Doris 是否连通、磁盘空间是否充足然后是同步主体节点最后是数据量校验节点。任务类型的选择要分场景。前置检查、数据量校验这类逻辑用 Shell 节点就够了里面写一些简单的 curl 或 MySQL 客户端命令。同步主体我们用 Python 节点因为要维护比较复杂的水位线逻辑和数据转换。DolphinScheduler 的 Python 任务支持 shell 里执行 python 脚本也可以直接用它的 Python API 定义任务我们用的是前者把脚本统一放在 AllData 管理的脚本仓库里版本可追溯。5.2 参数传递与定时调度DolphinScheduler 的参数体系一定要花时间理解清楚因为参数传错是任务跑飞最常见的元凶之一。参数分全局参数和工作流参数全局参数在整个工作流里所有任务节点都可以引用工作流参数只作用于当前节点。我们用得最多的还是日期参数比如$[yyyy-MM-dd]这种内置变量会在每个调度周期自动替换成当天日期。同步脚本并不直接把last_max_ts写死在代码里而是从 Doris 的状态表中读取所以 DolphinScheduler 的调度配置里不需要传太复杂的业务参数。定时调度直接配 cron 表达式比如20 * * * * ?表示每个整点过 20 分钟执行一次正好错开设备写入高峰。补数功能也是一定要留意的如果前一天凌晨的同步任务挂了等修完后直接在补数界面选日期范围重跑即可而不是手动改参数拼日期。5.3 任务失败企业微信告警配置告警这条我给新手一个建议一定要先把告警调通再上线同步任务。我们的告警规则是在 DolphinScheduler 的“告警组管理”里配置企业微信告警填入企业微信机器人的 Webhook 地址。工作流里每任务节点都可以指定告警策略我们设定为“失败时发送”并且失败重试次数设 3 次每次间隔 5 分钟。这么做的好处是偶发抖动会在 15 分钟内自动重试成功不需要打扰任何人如果 3 次重试都失败说明大概率是结构性故障比如 Doris 磁盘满了这时候企业微信告警才会发出。告警信息默认包含工作流实例 ID、任务名称、失败节点和日志链接。建议再配一个告警通知到 AllData 的统一入口这样值班人员在手机里收到告警后点进去就能看到这条同步链路的完整上下文而不是只看到一个孤零零的错误行。这个体验差距做过值班的人都能体会。5.4 资源与 Worker 分组避免任务互相踩踏同步任务看起来不重但物联网数据量一大读取、传输、写入都会消耗不少内存和网络带宽。如果 DolphinScheduler 的多个 worker 共享一台机器同步任务和别的任务抢资源的事情就会经常发生。我们当时的做法是给同步工作流单独建了一个 worker 分组指定 2 台专用的执行节点别的普通任务默认走另一个分组。这样即使同步任务计算量偶发暴涨也不会拖垮其他业务调度。Worker 分组同时也是一种故障域隔离。那 2 台同步节点如果硬件出问题只会影响同步工作流不会让整个调度平台瘫痪。从架构上看这个配置成本极低但稳定性收益很大强烈建议一上来就配上。6. 全流程踩坑排查从时序错乱到导入失败6.1 坑一时间窗口边界数据重复或丢失这个坑出现的频率非常高。我们的同步查询条件是ts last_max_ts AND ts end_time看似没有问题但只要同步任务在last_max_ts对应的时间点上有两条同一毫秒的测点记录就可能在两轮同步里都读到重复或都没读到丢失取决于查询时机。排查过程是这样的我们发现某天 Doris 的数据比 TDengine 多了一千多条比对每一行数据后重复的都是一些时间戳完全一样的记录。定位到问题后解决方案是把水位线从“时间戳”升级为“时间戳 批号”。每次从 TDengine 查询时除了记录最大时间戳还记录这一批数据的唯一标识比如设备 ID 列表的 hash下一轮查询不仅要求ts last_max_ts同时要求同一毫秒的数据只在下一轮里处理。简单说就是让时间边界带上数据指纹杜绝边界上的微妙重复。6.2 坑二Doris StreamLoad 报 label 冲突StreamLoad 导入时每次都要指定一个 labelDoris 靠 label 保证同一个导入任务的重试幂等性。但如果重试方向和业务预期不一致就会出问题比如第一次导入其实成功了但客户端因为网络超时认为失败脚本又重新发起导入Doris 检测到重复 label 直接返回Duplicate label这时候如果我们没正确解析这个错误就会把任务标记为失败启动重试形成死循环。我们的解法是在脚本里显式捕获 label 冲突异常去 Doris 查询该 label 对应的导入状态如果状态是 SUCCESS就认定当前批次已成功直接推进水位线不触发重试如果是其他状态再用新的 label 重新导入。同时在 DolphinScheduler 里把重试次数限制在 3 次以内防止异常情况下无限重试把 Doris 写入链路打满。6.3 坑三presto/doris 报 missing 的错误“presto doris 错误的 missing”这个热搜词背后其实是一个通用问题。如果业务方用 Presto 查询 Doris 数据偶尔会报出各种 missing 相关的错误。我当时排查下来主要原因多半是Presto 的元数据缓存和 Doris 的实际库表结构不一致。Doris 侧的库表刚创建或者字段刚变更而 Presto 的 information_schema 还没刷新Presto 在生成执行计划时就会认为某个列不存在报错 missing column。解决思路也很直接要么在 Presto 侧执行REFRESH SCHEMA要么给大数据量的小表建立数据库层面的CREATE TABLE IF NOT EXISTS之后再挂到 Presto。更重要的是建表和改表的流程要跟查询引擎的元数据刷新流程联动不能只在 Doris 执行 DDL 就结束。用 AllData 的好处就在这里我们把元数据变更做成一个带版本号的流程变更完成后自动触发下游引擎元数据刷新不再靠人肉记。6.4 坑四TDengine 查询超时和内存告警同步任务偶尔会在凌晨大查询时把 TDengine 的内存吃掉不少。后来看日志发现问题出在一条没有时间边界条件的 SQL 上调度刚启动时水位线表偶发读不到值脚本就 fallback 成全表扫描几亿行数据直接压过去完全绕过了所有优化。这个 bug 属于防御性编程没写到位。修复方法有两层第一层是在脚本里强制校验时间范围如果水位线缺失就默认从“当前时间往前推 24 小时”开始同步并要求最终输出本次查询的时间跨度超过阈值直接失败第二层是在 TDengine 侧设置查询内存限制和超时参数比如queryMemoryLimit和queryTimeout即使脚本写错了数据库也能兜底不会把整个节点搞挂。经过这次教训我每次写同步脚本都会把“兜底值”当成一等公民对待绝不假设上游状态一定正常。结尾说实话这套链路跑通不难难的是让它无人值守地稳定运行半年。我在这期间学到的最重要一条经验是同步任务的设计一定不要把逻辑都堆在一个脚本里而是要把“水位线管理”“数据导入”“状态记录”“异常上报”拆成清晰的模块各自独立可测。另一条经验是把平台能力用满DolphinScheduler 的告警和重试、Doris 的 label 幂等、TDengine 的超级表模型都是前人踩过坑沉淀下来的能力拿来就用比自己造轮子稳得多。如果你也在搭类似的物联网数据同步链路希望你读到这篇之后能少走我们当初那些弯路。