物联网项目做得多了你会发现采集其实是最省心的一步真正磨人的是把时序数据库里的海量数据搬到分析型数据库这一段路。最近我在一个工厂设备监测项目里用 AllData 平台统一管数据资产基于 DolphinScheduler 把 TDengine 里的设备指标按小时同步到 Doris前后跑了两个多月单日同步数据量过亿行。这篇文章我把整条链路拆开讲为什么选这套组合、表结构怎么设计、水位线怎么推进、Stream Load 有哪些坑、DolphinScheduler 任务怎么编排才稳给同样在做物联网数据同步的人一份可以直接抄作业的参考。先说项目背景。现场大概有 8000 多台设备每台设备每 5 秒上报一次数据包含温度、湿度、功率、电压等指标。TDengine 负责实时写入对外提供最近一周的时序查询Doris 负责跨设备的统计分析和长期报表。两者之间每天要搬运数亿行原始数据而且不能丢、不能重还要能回溯每一批数据从哪个时间点来、同步到哪个时间点。听起来不复杂但真搭起来接缝处全是细节。1. 项目背景与方案选型为什么是 TDengine 同步到 Doris1.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 数据链路全景先画清楚同步链路整体是这样的TDengine 超级表按设备和时间存储原始时序数据DolphinScheduler 每小时触发一次 Python 任务Python 脚本通过 TDengine 的 REST 接口按时间水位线增量查询数据查询结果转成 CSV 写入临时文件脚本调用 Doris 的 Stream Load 接口把 CSV 导入 Doris 明细表批次成功后更新 Doris 里的水位线记录表失败则走 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 平台里接入 DolphinSchedulerAllData 做的事情是把数据开发链路里的多个组件整合到一个入口管理。这里我主要用它来做三件事注册 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 侧表结构设计与抽取 SQLTDengine 建表时最核心的是区分标签和列。标签是用来过滤设备的列是真实测量值。我们的超级表结构大致是这样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 ) ENGINEOLAP 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 ) ENGINEOLAP 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_loadlabel是 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 重复提交确认是不是重试任务若需全新导入则换新 labelReach 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 幂等”这三个概念吃透天底下大部分离线同步任务都能照这个套路做。尤其是物联网数据量大、规律性强、时间属性明确其实是最适合用定时批量同步的。希望这篇流程能让你少踩几个坑同步链路一次跑顺。