☰
TDengine同步到Doris:物联网亿级数据离线同步全链路实践
2026/9/26 5:43:19 网站建设 项目流程

物联网项目做得多了,你会发现采集其实是最省心的一步,真正磨人的,是把时序数据库里的海量数据搬到分析型数据库这一段路。最近我在一个工厂设备监测项目里,用 AllData 平台统一管数据资产,基于 DolphinScheduler 把 TDengine 里的设备指标按小时同步到 Doris,前后跑了两个多月,单日同步数据量过亿行。这篇文章我把整条链路拆开讲:为什么选这套组合、表结构怎么设计、水位线怎么推进、Stream Load 有哪些坑、DolphinScheduler 任务怎么编排才稳,给同样在做物联网数据同步的人一份可以直接抄作业的参考。

先说项目背景。现场大概有 8000 多台设备,每台设备每 5 秒上报一次数据,包含温度、湿度、功率、电压等指标。TDengine 负责实时写入,对外提供最近一周的时序查询;Doris 负责跨设备的统计分析和长期报表。两者之间每天要搬运数亿行原始数据,而且不能丢、不能重,还要能回溯每一批数据从哪个时间点来、同步到哪个时间点。听起来不复杂,但真搭起来,接缝处全是细节。

1. 项目背景与方案选型:为什么是 TDengine 同步到 Doris

1.1 这类同步链路到底解决了什么问题

物联网数据的典型特征,是高频写入、按时间聚合、跨设备关联。数据源头的 TDengine 本质是时序数据库,它在写入速度和按时间的聚合查询上很擅长,但对复杂的多表 join、窗口分析、BI 报表这类场景并不顺手。而 Doris 是 MPP 分析型数据库,特别适合大宽表、高并发查询和物化视图,跑报表和临时 SQL 都比时序库舒服得多。

所以这套同步链路解决的不是“数据能不能导过去”的问题,而是:

  • 把实时写入层和分析查询层解耦,避免报表查询把写入链路拖垮;
  • 把原始时序数据按统一格式沉淀到 Doris 的明细表,供后续出指标、出报表;
  • 用一个可靠的调度平台管理同步任务,确保每小时一批、每批可重试、可监控,而不是靠 crontab 或手工脚本凑合。

如果数据量小,几千行可能无所谓,但一旦上了亿行级别,没有水位线管理和幂等机制,同步任务早晚会出事。

1.2 组件选型背后的逻辑

这套方案里每个组件都不是随便选的。

TDengine 承担采集层,最合适的理由就是它的超级表设计。一个 STable 能把所有设备统一成一张逻辑表,同时按设备打标签,查询单设备或单区域的数据很快。另一个理由是它自带保留策略,原始时序数据在 TDengine 里只保留最近一段时间,历史数据自动过期,省存储。

Doris 承担分析层,看中的是它的 Unique Key 模型和 Stream Load 导入能力。同一条时序记录重复导入时,Unique Key 可以按主键覆盖,这正好解决了调度任务重试导致重复数据的问题。Stream Load 走 HTTP 协议,支持 CSV 和 JSON,在调度任务里直接调用就行,不用额外部署一套传输组件。

DolphinScheduler 承担编排调度,因为它足够“重”——不是指它笨重,而是说它该有的功能都有。正排依赖、失败重试、超时控制、告警通道、工作流版本管理,这些在长期跑批任务里都是刚需。我们通过 AllData 平台把它和数据源元数据管理、数据质量管理统一收口,任务可以在一个入口里看完调度状态和血缘关系。

有些人会问,直接用 TDengine 的 taosAdapter 配合 Kafka 再进 Doris 是不是更实时?那确实更适合秒级同步场景。但我们的需求是小时级批量同步,Kafka 链路从运维成本和组件数量上看都是过度设计。

1.3 数据链路全景先画清楚

同步链路整体是这样的:

  1. TDengine 超级表按设备和时间存储原始时序数据;
  2. DolphinScheduler 每小时触发一次 Python 任务;
  3. Python 脚本通过 TDengine 的 REST 接口,按时间水位线增量查询数据;
  4. 查询结果转成 CSV 写入临时文件;
  5. 脚本调用 Doris 的 Stream Load 接口,把 CSV 导入 Doris 明细表;
  6. 批次成功后,更新 Doris 里的水位线记录表;
  7. 失败则走 DolphinScheduler 重试机制,整个批次重新执行。

这套设计里没有额外引入 DataX 或 SeaTunnel,不是因为它们不行,而是当前场景下 Python 脚本直连两端接口,链路最短、排障也最直接。后面如果数据量再翻倍,可以在中间加文件暂存或换 SeaTunnel,但那是后话。

2. 环境准备:TDengine、Doris、DolphinScheduler 的落地细节

2.1 TDengine 部署与保留策略

TDengine 社区版安装本身不复杂,但有几个关键点必须说。

第一是磁盘规划。时序数据压缩率不低,但写入量大的时候,数据文件和 WAL 文件的 IO 竞争很厉害。我在部署时把数据目录和 WAL 目录分到了不同磁盘,实测下来写入毛刺明显减少。

第二是保留策略。TDengine 的保留策略在创建数据库时通过KEEP设置,单位是天。我们这个项目只需要保留 7 天原始数据,所以建库语句大概是这样:

CREATE DATABASE iot_db BUFFER 256 CACHEMODEL none WAL 20 KEEP 7 STTAGR 100;

STTAGR是阈值,超过这个行数的子表数据会被合并到列存中,这个参数调大可以减少小文件数量,但查询最新数据时延迟会增加。我调过几次,最终觉得默认值附近最省心。

第三是账号权限。虽然 TDengine 对权限分离做得不重,但建议单独建一个只读账号给同步任务用,别用 root 去跑。这样排查问题的时候,至少能区分是查询端的问题还是写入端的问题。

2.2 Doris 集群部署要点

Doris 部署主要分 FE 和 BE 两类节点。FE 负责元数据和查询解析,BE 负责数据存储和计算。

我们这个规模用了 1 个 FE 加 3 个 BE 的部署,BE 每台机器给 64GB 内存。按 Doris 的惯例,BE 内存主要是给执行查询和导入用的,如果机器内存不够,Stream Load 大批量导入时很容易触发内存不足。

这里有一个经常被忽略的点:Doris 的 Stream Load 是通过 FE 转发或重定向到 BE 的,默认访问端口是 8030。部署完最好先验证一下这个端口的连通性,很多调度脚本导入失败,根本不是导入 SQL 的问题,而是防火墙把 BE 的 8030 端口挡住了。

我遇到过最离谱的情况是 FE 能通、BE 连不上,Doris 返回一个“redirect”响应,然后 Stream Load 超时。所以在环境准备阶段,最好把 FE 的 8030 和 BE 的 8040 端口连通性都测一遍。

2.3 AllData 平台里接入 DolphinScheduler

AllData 做的事情是把数据开发链路里的多个组件整合到一个入口管理。这里我主要用它来做三件事:

  • 注册 DolphinScheduler 作为统一的调度引擎;
  • 把 TDengine 和 Doris 的连接信息做成可复用的数据源;
  • 把调度任务和元数据血缘挂上,后续看一条数据从哪来到哪去一目了然。

在 AllData 中接入 DolphinScheduler 的步骤相对机械:组件管理里选择 DolphinScheduler 服务地址,填入管理员账号,然后它会把工作流、任务实例同步过来。接入后可以直接在 AllData 页面上创建项目、绑定数据源、发布工作流。

数据源层面的一个建议是:TDengine 数据源类型如果选项里找不到,就用 REST API 类型代替,连接串写成 TDengine 的 REST 地址。Doris 如果系统里没有专门的 Doris 类型,可以选 MySQL 协议连 FE 的 9030 端口,因为 Doris 兼容 MySQL 协议。

3. 数据同步全流程实现:从 TDengine 抽取到 Doris 装载

3.1 TDengine 侧表结构设计与抽取 SQL

TDengine 建表时最核心的是区分标签和列。标签是用来过滤设备的,列是真实测量值。

我们的超级表结构大致是这样:

CREATE STABLE iot_db.device_data ( ts TIMESTAMP, temperature FLOAT, humidity FLOAT, power FLOAT, voltage INT ) TAGS ( device_id VARCHAR(32), location VARCHAR(64) );

设备在接入时会单独创建子表,子表名字通常是设备编码。这样设计的好处是,TDengine 对每个子表的写入是天然隔离的,查询时按device_id过滤会走标签索引,速度很快。

抽取的 SQL 一定不要写成只按时间范围扫全表,而要带上标签条件。如果你是全平台同步,那至少按区域或设备分组并行抽取,避免一个查询把 TDengine 的内存打爆。我的做法是先把设备清单按区域拆成 4 组,每组对应调度工作流里的一个并行分支,每个分支查一组设备。这样单次查询的行数可控,失败重试的影响面也小。

抽取 SQL 的典型形态是:

SELECT ts, device_id, location, temperature, humidity, power, voltage FROM iot_db.device_data WHERE device_id IN ('device_a', 'device_b', ...) AND ts > '2025-06-01 00:00:00' AND ts <= '2025-06-01 01:00:00' ORDER BY ts ASC;

ORDER BY ts是为了保证写入 Doris 后按时间有序,方便 Doris 的存储层做更好的索引剪裁。

3.2 增量水位线怎么设计

增量同步最怕的就是“不知道上次同步到哪了”。如果每次都全量同步,数据量一大,任务时间就会越过调度周期;如果不记录水位线,漏数据都没法发现。

我们建了一张水位线表存在 Doris 里,结构很简单:

CREATE TABLE sync_watermark ( src_table VARCHAR(128), sync_target VARCHAR(128), last_watermark DATETIME, update_time DATETIME ) ENGINE=OLAP UNIQUE KEY(src_table, sync_target) DISTRIBUTED BY HASH(src_table) BUCKETS 3;

同步任务启动时先查这张表,拿到上次同步到的最大时间;任务成功后再把本次批次的最大时间更新回去。

这里有一个关键点:水位线更新和 Stream Load 导入一定要放在同一个任务的后续步骤里,不能并行动作。否则可能出现数据还没导完,水位线已经推进了,下一次同步直接跳过一批数据。等到项目稳定后,还可以把水位线更新放到独立任务节点,通过 DolphinScheduler 的依赖关系控制先后顺序,这样重试时不会误改水位线。

3.3 临时数据落地与清洗

抽取出来的原始数据不能直接砸进 Doris,至少要过一遍清洗。

我习惯让 Python 脚本把 TDengine 的查询结果先转成 CSV 字符串,然后写到本地临时目录。这个临时目录最好和 DolphinScheduler 的工作目录分开,而且要预留足够空间。每小时几千万行的 CSV,差不多要占 1GB 到 2GB 空间,如果磁盘满了,导入任务会报错,而且这种错误特别难排查,因为报错信息指向的是文件写入失败,而不是导入失败。

清洗主要做三件事:

  • 把 TDengine 返回的时间格式统一成 Doris 能识别的yyyy-MM-dd HH:mm:ss;
  • 把NULL或空值替换成 Doris 可以处理的默认值;
  • 把字符串字段首尾的空格去掉,避免 Doris 里出现看起来一样但实际不等的脏数据。

清洗逻辑不复杂,但要在脚本里形成固定流程,不要指望数据源侧全干净。

3.4 Doris 建表与 Stream Load 写入

Doris 建表需要考虑查询场景。我们分析层主要按设备和时间查明细,所以明细表的模型选了 Unique Key,这是为了配合同步重试时去重。

目标表大致是:

CREATE TABLE doris_db.device_data ( ts DATETIME, device_id VARCHAR(32), location VARCHAR(64), temperature DOUBLE, humidity DOUBLE, power DOUBLE, voltage INT ) ENGINE=OLAP UNIQUE KEY(ts, device_id, location) DISTRIBUTED BY HASH(device_id) BUCKETS 16 PROPERTIES ( "replication_num" = "3", "compression" = "ZSTD" );

分桶键选择device_id,是因为查询基本都会带上设备条件。如果查询按区域过滤多,也可以把区域放到分桶键里,但键太多会导致数据分布不均,我建议尽量少用多列分桶。

导入端我们用 Stream Load,这是 Doris 导入小批量文件最高效的方式。核心调用如下:

curl --location-trusted \ -u doris_user:doris_pass \ -H "label:iot_sync_202506010100_001" \ -H "column_separator:," \ -H "format:csv" \ -H "columns:ts,device_id,location,temperature,humidity,power,voltage" \ -H "strict_mode:false" \ -T /tmp/sync/device_data_202506010100.csv \ http://doris-fe:8030/api/doris_db/device_data/_stream_load

label是 Stream Load 幂等的关键。同一个 label 重复提交,Doris 会直接返回第一次提交的结果,而不会重复导入。重试机制全靠它撑着。

有一个参数我专门说一下:strict_mode。默认情况下它是 false,导入时对字段类型不匹配比较宽容。如果你希望脏数据直接报错而不是放过,才需要开 true。我们这里选择 false,是因为偶尔会有设备上报异常值,比如电压变成负数,严格模式会把整批数据拦截,影响任务稳定性。

导入完成后,脚本需要解析 Doris 返回的 JSON,重点看Status字段。状态为Success才能推进水位线,Fail或Label Already Exists都要当作异常处理。

3.5 DolphinScheduler 工作流编排与发布

DolphinScheduler 里整个同步工作流分三层:

第一层是参数定义。我习惯把 TDengine 的 REST 地址、Doris 的 FE 地址、批次大小、设备分组数都定义成工作流全局参数,这样改环境时不用逐个改任务。

第二层是任务节点。每个区域分组对应一个 Python 任务,四个任务并行执行,互不依赖。Python 脚本上传到 DolphinScheduler 的资源中心,然后任务节点直接引用资源文件。这样脚本改动可以在资源中心里新版本,不会污染正在执行的任务。

第三层是调度周期。DolphinScheduler 的调度配置里选择 cron 表达式,每小时在整点过 5 分钟开始跑,也就是0 5 * * * ? *。这个不是随便拍的,错开整点的目的是避免和数据源上报高峰撞车,把同步任务对 TDengine 查询性能的影响降到最低。

工作流发布后,我建议先跑几个手动实例,确认 batch 标签唯一性、水位线更新正常,再打开周期调度。不然一上来就自动跑,出问题都不知道哪一批写坏了。

4. 常见问题与排查技巧实录

4.1 时区与时间格式不一致

这个坑几乎必踩。TDengine 的 REST 接口返回的时间默认是 ISO8601 格式,带有时区后缀;而 Doris 的 DATETIME 字段要求的是纯日期时间字符串。直接拿 TDengine 的原始值导入,Doris 会报时间格式错误。

我们的解决方案是在 Python 脚本里统一做一次时间转换,把 TDengine 返回的字符串解析成 Python 的datetime,再格式化成%Y-%m-%d %H:%M:%S输出。时区方面,所有集群都统一用 Asia/Shanghai,避免跨时区导致的边界错位。

另一个细节是 TDengine 的 SQL 查询里,时间条件也要注意带不带毫秒。两次同步批次之间,如果上一批的结束时间和下一批的起始时间重叠,在 Doris 的 Unique Key 下覆盖没问题,但查询时可能出现重复计算的假象,所以时间条件建议用半开半闭区间。

4.2 Stream Load 报错的中英文对照与处理

Stream Load 的报错信息比较直白,但网络上一搜答案特别散,我这里整理几个高频的:

报错现象实际原因处理方式
Table xxx has xxx replicas, but only 1 aliveBE 节点宕机或网络不通先查 BE 状态,再查防火墙端口,最后查磁盘空间
The specific label has been used相同 label 重复提交确认是不是重试任务;若需全新导入,则换新 label
Reach limit of connectionsFE 连接数满了调大 FE 的qe_max_connection或检查是否有连接泄漏
Timeout数据量过大或导入时间超限缩小每个分组的设备数,或调大 Stream Load 的超时时间

我建议在 Python 脚本里把 Doris 返回的完整 JSON 打印到日志里,不只打印状态码。很多时候问题出在ErrorURL字段指向的具体错误文件,不打开它排查不出来。

4.3 任务重试导致重复数据和脏数据

任务重试造成的重复数据,理论上会被 Doris 的 Unique Key 覆盖,但前提是你的主键设计够完整。如果主键只用了ts而没带device_id,两个设备同一时间点的数据就会互相覆盖,这不是语义上的“去重”,而是数据丢失。

我们最终把主键定为ts + device_id + location,其中 location 是冗余的。加它的原因是有些设备迁移过位置,同一时间点上同一 device_id 可能对应不同 location,如果主键不带 location,历史数据会被新位置覆盖,查询时历史归因就会错。

如果发现导入的批次有脏数据想回退,光靠 Unique Key 不行。Doris 的 Stream Load 本身不支持按 label 回滚,只能通过删除范围内的数据再重新导入修复。所以我在工作流里专门留了一个“重刷最近 N 小时”的备用任务,遇到数据质量问题时先删对应时间窗,再重新同步。

4.4 资源竞争与链路性能优化

同步任务跑起来后,最大的性能瓶颈其实不在 Doris,而在 TDengine 的查询端。如果你每小时的批次全压在几分钟内发起查询,TDengine 的内存会飙升,WAL 写入也会被拖慢。

我做了三个优化:

  • 错峰:不同设备组的查询时间错开,而不是同一秒全部启动;
  • 限流:Python 脚本里对 TDengine REST 接口做了简单的请求间隔控制,避免短时间并发过高;
  • 分页:单次查询如果超过 500 万行,就按时间窗再切小,分页拉取,防止 TDengine OOM。

Doris 侧的优化主要是分桶数。如果 BE 有 3 台,分区桶数设为 3 的整数倍比较均衡。我们一开始 BUCKETS 设 8,后来发现个别桶的数据量明显偏大,改成 16 才匀了一些。分桶数不是越大越好,过大会产生很多小文件,合并压力反而变大。

4.5 调度实例卡死和假死排查

DolphinScheduler 跑久了,偶尔会出现任务实例一直显示“运行中”,但实际脚本已经退出的情况。多数时候是 Python 脚本里的某个连接没有关闭,进程挂住。排查方法很简单:进到 DolphinScheduler 的工作目录,看任务对应的日志尾部,如果日志停在某个网络请求后不再输出,大概率是连接等待超时。

我的经验是给 Python 脚本里所有 HTTP 请求都加上超时参数,并且设置重试次数。否则一次网络抖动,任务挂几小时,后续批次全部排队。DolphinScheduler 的失败重试次数默认是 0,把重试次数设置成 2 到 3 次,重试间隔 5 分钟,基本能 cover 大部分网络瞬时故障。

把所有工作流跑稳之后,这个项目最大的收获就是:同步链路不是堆组件,而是把数据边界、幂等机制和排障手段设计好。每次 Doris 里查出来的数据,我都能顺着 label 和水位线倒推它是从 TDengine 的哪个时间窗口来的,出了问题也知道该去哪里修数据。

这里也顺带提醒一句,别贪多求新,把 DataX、SeaTunnel、Flink CDC 全塞进来。链路越长,能出问题的接缝越多。我后来在另一条轻量级数据同步需求里也试过用 DolphinScheduler 直接调度 Shell 脚本加 DataX,但只要数据量没到百万亿级别,简单方案永远更好维护。

这几个月跑下来,我个人体会最深的一点是:把“批次”“水位线”“label 幂等”这三个概念吃透,天底下大部分离线同步任务都能照这个套路做。尤其是物联网数据,量大、规律性强、时间属性明确,其实是最适合用定时批量同步的。希望这篇流程能让你少踩几个坑,同步链路一次跑顺。

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

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

立即咨询