ARTICLE DETAIL

资讯详情

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

InfluxDB数据同步到Doris:基于SeaTunnel的完整落地实践

InfluxDB数据同步到Doris:基于SeaTunnel的完整落地实践 去年我们在给监控平台做升级的时候遇到一个很典型的需求业务指标、服务器监控数据全部落在InfluxDB里平时查短期趋势没毛病但一张报表如果要关联业务订单表、按天做多维聚合InfluxDB就明显吃力了。后来我们把分析侧的查询整体迁移到Doris中间就面临一个绕不开的问题数据怎么从InfluxDB稳定地同步到Doris。折腾了一圈最后选定的方案是用SeaTunnel来做同步管道。这套组合我们已经跑了小半年单日几十亿点位的时序数据同步到Doris稳定性和性能都达到了生产要求。这篇文章就把整个落地过程、核心配置参数和踩过的坑整理出来给同样在搞时序数据实时数仓的兄弟做个参考。文章里你会看到为什么选SeaTunnel而不是DataX或者自研脚本InfluxDB的时间戳和tag在Doris里怎么设计表结构以及一套可以直接抄作业的同步任务配置。不管你是刚接触这三样东西还是已经在做数据接入但被各种小问题卡住这篇都值得看完。1. 数据管道为什么是InfluxDB加Doris中间又为什么是SeaTunnel1.1 InfluxDB解决的是写入和短时查询不是分析先捋一下角色定位。InfluxDB是时序数据库它最擅长的场景是海量监控指标、IoT传感器数据的持续写入和短时间范围内的趋势查询。我们当时所有服务器的CPU、内存、磁盘IO、接口调用量都写进了InfluxDB单机每秒写入几万点查询最近一小时的数据响应都是毫秒级。但一旦查询范围变大比如你要从几亿条时序数据里按天聚合出每台机器的平均使用率再和业务库的订单表做关联InfluxDB就很吃力了。InfluxQL的分析能力偏弱多表join、复杂子查询基本没法做而且大范围扫描的查询会把内存打满影响正在进行的写入。再一个BI工具连InfluxDB做报表体验也不好大部分报表工具对时序数据库的支持都停留在能连上、能出图的程度离真正的自助分析差很远。1.2 Doris补的是大规模多维分析这块短板Doris是MPP架构的实时分析型数据库列式存储加上向量化执行再做高并发的多维聚合、关联查询非常顺手。我们的报表、即席查询、大屏数据全部迁移到Doris之后原来要跑几十秒的聚合SQL现在基本一秒内出结果。有人会问那直接用StarRocks不也行吗确实Doris和StarRocks同源功能和性能很接近。我们最后选Doris是综合考虑了社区活跃度、版本迭代节奏和团队已有技术栈的维护成本。如果你已经在用StarRocks下面这套SeaTunnel同步方案的核心思路同样适用把Doris Sink换成语法兼容的下游即可。1.3 SeaTunnel凭什么能当这个搬运工先说说最原始的方案自己写程序从InfluxDB查数据解析完再通过Doris的Stream Load接口写入。这个方案看着灵活实际做起来很麻烦。你要处理连接管理、失败重试、批量写入、类型转换、断点续传这些问题写出来的代码基本是一坨只属于你们团队、别人不敢动的“祖传逻辑”。SeaTunnel是一个开源的数据集成工具核心思路是把你需要的数据管道用一套配置文件描述出来source数据源、transform转换、sink目标端三段式配置写完直接跑。它自己带了一个Zeta引擎不用再装Spark或者Flink单机也好集群也好一条命令就能起任务。对比DataXSeaTunnel的实时性更好Sink插件的生态也更丰富。对比自己写脚本SeaTunnel把分布式调度、失败重试、checkpoint这些机制都内置了一个同步任务写下来配置文件一百行以内搞定维护成本低太多。2. 开始前的环境准备InfluxDB、Doris、SeaTunnel三板斧2.1 InfluxDB侧的准备版本确认、数据写入验证、磁盘检查我这里以InfluxDB 1.8为例这是目前生产环境里最稳的1.x版本API稳定文档也多。用二进制包或者Docker都行Docker一键起服务最省事docker run -d --name influxdb \ -p 8086:8086 \ -v /data/influxdb:/var/lib/influxdb \ influxdb:1.8启动之后用客户端建一个库或者直接写几条测试数据进去确保InfluxDB本身是通的# 进入容器 docker exec -it influxdb bash # 建立monitoring库 influx -execute CREATE DATABASE monitoring # 写入一条CPU监控数据tag是host和regionfield是usage和load influx -database monitoring -execute INSERT cpu,hostserver01,regioncn usage42.5,load0.8建好库、写入数据之后用查询验证一下influx -database monitoring -execute SELECT * FROM cpu会看到类似下面的结果time是RFC3339格式的时间字符串后面跟着tag、field字段time host load region usage ---- ---- ---- ------ ----- 2024-01-01T00:00:00Z server01 0.8 cn 42.5这里要提醒一句如果你用的是InfluxDB 2.xSeaTunnel连接的时候建议走它的1.x兼容API在连接串和账号密码上沿用1.x的方式比直接用Flux查询要省事得多。具体做法是在influxdb.conf里把[http]下的flux-enabled设为true然后请求地址用/api/v1前缀。另外排查问题的时候强烈建议装一个InfluxDB Studio桌面客户端可视化地看看数据长什么样。很多时候配置写不对就是因为对InfluxDB里字段的真实类型和值没概念。用Studio先跑一遍查询确认返回的列名、类型再往SeaTunnel配置里填能省掉大量来回试错的时间。还有一个坑InfluxDB如果报engine: error writing wal entry: write /var/lib/influxdb/wal/krakend/autogen这种WAL写错误大概率是磁盘满或者WAL目录权限不对。先把磁盘空间清出来确认目录的属组是influxdb用户再重启服务别一上来就怀疑同步工具。2.2 Doris侧的准备集群部署和明细表设计Doris的部署本身不复杂FE负责元数据和查询解析BE负责数据存储和计算。下载官方二进制包解压之后先起FE再起BE# 起FE cd apache-doris/fe/bin ./start_fe.sh --daemon # 起BE cd apache-doris/be/bin ./start_be.sh --daemon然后用MySQL协议连上FE把BE节点加进去-- 用mysql客户端端口是FE的query_port默认9030 mysql -h 127.0.0.1 -P 9030 -uroot -- 添加BE节点 ALTER SYSTEM ADD BACKEND 127.0.0.1:9050;加完之后可以通过SHOW BACKENDS确认BE状态是Alive。集群管理层面如果机器数量多建议再装一个Doris Manager用来监控节点状态、做告警、管理用户权限比纯命令行省心很多。接下来是建表。InfluxDB的一条数据本质上由三部分组成时间戳、tag、field。同步到Doris之后的表结构设计首要原则是尽量保留源数据的“明细程度”方便后续做任意维度的分析。我推荐用Duplicate模型把所有tag字段放到Key列field字段放在Value列CREATE TABLE monitoring.cpu_usage ( ts DATETIME, host VARCHAR(64), region VARCHAR(32), usage DOUBLE, load DOUBLE ) DUPLICATE KEY(ts, host, region) DISTRIBUTED BY HASH(host) BUCKETS 10 PROPERTIES ( replication_num 1 );为什么用Duplicate而不是Aggregate或Unique因为监控数据本身就是“只追加、不更新”的明细数据后续分析需要的是最大粒度的原始值。如果建表时就把usage定义成Aggregate模型下的SUM或AVG后面的分析场景就被限制死了。时间戳和host、region放在Key列是为了让相同维度的数据天然落在一起按时间范围扫描的时候性能更好。2.3 SeaTunnel侧的准备安装和连接器SeaTunnel安装更简单下载发行包解压就能用。需要注意版本和JDK的对应关系2.3.x系列要求JDK8以上。我用的版本是2.3.8界面没有太多花哨的东西核心就是启动脚本加配置目录。默认发行包不会带上所有连接器需要自己装。执行安装脚本把InfluxDB和Doris的connector装进去# 进入SeaTunnel目录 cd /opt/seatunnel # 安装influxdb和doris的connector2.3.5之后推荐这样按需安装 sh bin/install-plugin.sh --name connector-influxdb --name connector-doris装完之后在connectors/目录下能看到对应的jar包比如connector-influxdb-2.3.8.jar和connector-doris-2.3.8.jar。如果启动任务时报找不到插件类八成是这一步没做或者jar包的版本和主程序不一致。3. 核心实操编写并跑通第一个同步任务3.1 配置文件拆解env、source、transform、sinkSeaTunnel的任务就是一个配置文件核心由四个块组成。我把我们生产环境用的配置简化了一下拿来直接说明每一段的作用。假设我要同步monitoring库里cpu这个measurement目标表是Doris里的monitoring.cpu_usage同步最近一小时的数据。完整配置如下env { parallelism 2 job.mode BATCH } source { InfluxDB { url http://192.168.1.10:8086 database monitoring username root password root sql SELECT time, host, region, usage, load FROM cpu WHERE time now() - 1h schema { fields { time STRING host STRING region STRING usage DOUBLE load DOUBLE } } partition_column host partition_num 2 fetch_size 1000 } } transform { # 这里先留空后面讲时间格式转换的时候再展开 } sink { Doris { fenodes 192.168.1.20:8030 username root password table.identifier monitoring.cpu_usage column_names [ts, host, region, usage, load] doris.config { format json read_json_by_line true } } }先看env块。parallelism是任务并行度我设成2表示数据会被拆成2个分片并行读取和写入。job.mode BATCH表示这是一个批处理任务跑完就结束不会常驻。如果是想持续同步可以换成STREAMING模式但配合调度系统做定时批处理在大多数场景下更可控、更好排查问题。再看source块。url、database、username、password是InfluxDB的连接信息。sql就是一次普通的InfluxQL查询这里有一个非常关键的点sql里查出来的列必须和下面schema定义的字段一一对应。SeaTunnel会按照schema里的字段顺序去解析查询结果多了少了都会报错。schema里的字段类型需要和InfluxDB返回值的真实类型匹配。time字段在InfluxDB的查询结果里是字符串所以在schema里定义成STRING。host、region本身是tag默认就是字符串类型。usage和load是field这里定义成DOUBLE对应Doris里的DOUBLE列。partition_column host和partition_num 2表示按host字段做切分让两个并行度分别处理不同host的数据。这个设计能大幅度提升大批量同步的吞吐后面讲调优的时候再细说。sink块里fenodes是Doris FE的HTTP端口默认是8030SeaTunnel会通过这个地址发Stream Load请求。table.identifier就是目标库表名。column_names是写入目标表的列名列表顺序要和数据行保持一致。这里我把InfluxDB查询出来的time映射到了Doris表的ts列其余字段名保持一致。doris.config里面塞的是Doris Stream Load的参数这里用JSON格式写入同时开启按行读取JSON。Stream Load是Doris提供的一种高效批量导入方式走HTTP协议SeaTunnel底层就是用它把数据灌进Doris的。3.2 时间格式转换同步任务最容易被坑的地方配置写好之后先别急着跑有一个问题必须处理时间字段。InfluxDB查询出来的time是RFC3339格式长这样2024-01-01T00:00:00.000000000Z。而Doris里的ts列是DATETIME类型接受的格式是2024-01-01 00:00:00或者带毫秒的2024-01-01 00:00:00.000。如果不做转换直接塞进去Doris会解析不了Stream Load直接报错。处理办法是在InfluxDB的SQL里就把时间格式转好用InfluxQL的time_format函数。把source块里的sql改成SELECT time_format(time, yyyy-MM-dd HH:mm:ss) AS ts, host, region, usage, load FROM cpu WHERE time now() - 1h这样查出来的time字段就变成了2024-01-01 00:00:00正好是Doris支持的格式。schema里依然保持STRINGDoris的DATETIME可以自动解析这个格式的字符串。如果你的时间字段不想在源头处理而是在SeaTunnel的transform块里做也可以写一个Copy或者Replace的转换。但我个人建议能压在SQL里就压在SQL里少一层处理就少一个故障点InfluxQL本身自带的函数已经够用了。3.3 首次运行提交任务、看日志、验证结果配置改好之后怎么提交SeaTunnel提供了Zeta引擎的提交方式在安装目录执行bin/seatunnel.sh -c config/influxdb_doris.conf -t zeta-c指定配置文件路径-t指定引擎类型zeta就是SeaTunnel自带的引擎。任务启动后日志会打到控制台也可以配置写到日志文件里。正常运行的时候日志里会出现类似下面这样的关键信息INFO Job Statistic: { ReadCount: 50000, WriteCount: 50000, TotalReadTime: 123456789, ... }ReadCount和WriteCount分别表示读取和写入的行数两者相等就说明这批数据全部成功写入。如果WriteCount比ReadCount少说明有部分数据写入失败SeaTunnel会抛出异常把具体的失败原因打到日志里。任务跑完之后去Doris里验证一下数据SELECT COUNT(*) FROM monitoring.cpu_usage; SELECT * FROM monitoring.cpu_usage ORDER BY ts DESC LIMIT 10;能看到数据而且ts格式正常、字段值和InfluxDB里的一致第一个同步任务就算彻底跑通了。4. 生产环境落地增量同步、并发调优和常见性能问题4.1 增量同步思路定时任务加时间窗口过滤实际生产场景里全量同步很少跑跑一次之后基本都是增量。增量同步的核心思路很简单每次只查从上一次同步点到现在的新增数据。最简单可靠的时间窗口方案是用定时调度让任务周期性执行每次都同步最近10分钟或者最近1分钟的数据。比如你用crontab每5分钟跑一次SQL里就直接写SELECT time, host, region, usage, load FROM cpu WHERE time now() - 10m这样的做法对InfluxDB的查询压力很小而且天然具备断点续跑能力——即使某次任务失败下个周期再跑的时候只要时间区间覆盖住之前的数据就能补回来。这里注意一下InfluxQL里的时间单位m表示分钟h表示小时d表示天写错单位会直接导致查不到数据。如果对断点精度有更高要求比如每1分钟跑一次那就在每台执行调度的机器上维护一个游标文件记录上次同步的最大时间戳。任务启动时读取这个游标查询大于等于游标值的数据任务结束后用这次查到的最大时间戳更新游标。这个方法我建议大家在批处理方案稳定之后再去优化前期用固定时间窗口完全够用。4.2 并行度和分片参数怎么调一开始我们同步全量数据几亿条数据要跑几个小时后来发现问题是并行度完全没生效。问题就出在partition_column上。SeaTunnel的InfluxDB Source支持按某个字段对查询结果做分片并行度的上限由partition_num决定。我当时在host字段上做了分区单条SQL变成了多个并行的查询写入Doris的吞吐瞬间就上去了。调并行度的时候要有个度。并行度太高InfluxDB那边查询并发太大可能会把时序库的CPU打满影响线上指标的实时写入。我这边单节点InfluxDBparallelism设成4partition_num设成4同步速度已经非常可观。如果InfluxDB是集群部署可以适当加大。Doris侧也有一个参数值得调。sink块里的batch_size控制着单个批次写入的数据量默认配置下如果单批太小写到Doris的请求数就会很多反而影响吞吐。我习惯把它调节到对应每批几万行的水平具体值取决于单行数据的大小。类似调整还有sink.max-retries控制写入失败时的重试次数建议设置成3配合重试间隔能够消化掉一部分短暂的网络抖动。4.3 内存和GC问题SeaTunnel任务OOM怎么办SeaTunnel跑大批量同步的时候如果任务突然崩掉日志里报OOM内存溢出首先检查两处一是SeaTunnel进程的JVM堆内存二是fetch_size。SeaTunnel的启动脚本默认堆内存可能只有1GB同步数据量大、并行度高的时候完全不够。解压目录下的bin/seatunnel.sh可以调整启动参数把-Xmx改成4GB或者更高。改完JVM内存再检查source块里的fetch_size这个值控制每次从InfluxDB取回多少条记录。我遇到过一次问题fetch_size设成10000配合并行度4一次任务就要缓冲四万行数据内存压力一下子就上来了。后面调成1000问题就消失了。这个调参的优先级要注意先动分区和并行度再看单批次大小最后才考虑加内存。无脑加内存只会让问题出现的阈值变高但GC停顿还是会在数据量大的时候拖慢任务。4.4 同步延迟优化关注Doris的Stream Load而不是SeaTunnel本身从SeaTunnel到Doris的写入性能瓶颈往往不在SeaTunnel而在Doris侧Stream Load的处理效率。如果发现写入延迟很高或者Doris的BE节点CPU飙升去Doris Manager上看一看Stream Load有没有堆积同时确认目标表的bucket数是否合理。bucket数也就是DISTRIBUTED BY HASH(host) BUCKETS 10里的10这个值和集群规模、数据量有关。数据量特别大的表bucket太少会导致单个BE分片数据倾斜写入速度上不去。一般建议单个bucket的数据量控制在100MB到1GB之间可以先按这个标准估算不够再加。加bucket数的操作需要表重建所以最好在表设计阶段就规划好避免后面来回折腾。5. 常见问题与排查技巧实录5.1 时间格式报错列解析失败现象任务跑起来Doris侧写入失败日志提示无法解析ts列的值或者转换类型出错。排查思路先用InfluxDB Studio跑一遍同步SQL看看time返回的具体格式。如果还是RFC3339带T和Z的格式说明SQL里的time_format函数没生效检查是不是函数名写错了。InfluxQL里时间是UTCDoris如果存储的是本地时间要注意时区偏移可以把Doris的time_zone参数和业务要求的时区保持一致。5.2 InfluxDB认证没开导致连接失败现象SeaTunnel日志报401 Unauthorized或者连接超时。排查思路确认InfluxDB的[http]配置里auth-enabled是true还是false。如果没开启认证SeaTunnel里的username和password随便填也能连上但其实建议在InfluxDB侧开启认证特别是内网多个服务共享一个InfluxDB的时候避免误删数据。5.3 Doris Stream Load报错too many open files现象任务跑了一段时间后日志大量报too many open files。排查思路这是文件句柄数不够通常因为任务持续建立HTTP连接而系统默认限制比较低。在启动SeaTunnel的机器上调整进程的nofile限制ulimit -n 65535同时如果Doris的BE节点也报类似的错去BE的配置文件里把max_open_files调大重启BE生效。5.4 连接器版本不一致导致加载失败现象提交任务时直接报ClassNotFoundException或者提示找不到InfluxDB Source插件。排查思路SeaTunnel对连接器的版本要求比较严格主程序和connector的版本必须一致。比如你用的是2.3.8的主程序插件jar包的版本也应该是2.3.8。另外有些老插件需要放在connectors目录下新版本推荐把插件jar包放在connectors/而不是lib目录下位置不同也会导致加载失败。如果不确定直接在解压目录全局搜一下connector的jar包在哪个位置对照官方文档调整。5.5 Doris CCR Syncer和外部数据接入工具是两个层面的东西见过不少人在做数据接入的时候搜到Doris CCR Syncer以为这也是一个同步工具。这里强调一下Doris CCR Syncer是Doris集群和Doris集群之间做数据同步的工具主要用于容灾、读写分离、机房级数据复制解决的是Doris内部的副本和数据一致性问题。而我们这里做的是“外部数据源InfluxDB到Doris”的接入完全不是一个层面的事。如果你需要的是业务侧的数据实时接入着眼点应该是DataX、SeaTunnel、Flink CDC这类数据集成工具别在CCR上浪费太多时间。5.6 常见问题速查表问题现象可能原因解决手段任务报错提示time字段解析失败时间格式不是Doris支持的DATETIME格式在InfluxQL里使用time_format函数按yyyy-MM-dd HH:mm:ss格式化SeaTunnel连接InfluxDB报401认证信息错误或鉴权未开启核对InfluxDB的auth-enabled配置和账号密码大量数据同步时任务OOMJVM堆内存不足或fetch_size过大调大SeaTunnel的-Xmx调小fetch_size到1000左右Doris侧报too many open files系统文件句柄数不足修改启动用户和BE进程的nofile限制写入行数和读取行数不一致存在部分脏数据导致Doris拒绝写入查看完整异常日志定位具体是哪几行数据格式有问题同步速度特别慢并行度未生效或partition_num设置不合理检查是否配置了partition_column和partition_num确认parallelism大于1日志提示找不到InfluxDB Source类connector插件未安装或版本不匹配重新执行install-plugin.sh确认插件jar包版本与主程序一致6. 一点个人体会和后续可以做的事这套InfluxDB到Doris的同步管道跑通之后最直观的感受是分析侧再也不受时序库的查询能力限制了。原来只能在一个时序库里看单条曲线的数据现在到了Doris里可以和业务数据做关联、做透视、做各种魔改报表。整个同步过程的维护成本也低SeaTunnel的配置化设计让新接一个数据源变得非常简单后期扩展基本就是改一个配置文件的事。最后再分享一个小技巧如果你对数据一致性要求比较高不想出现同步中断后重复插入导致的数据重复可以在Doris表里加一个SET类型或者使用Unique模型的Key列把时间戳和所有维度的tag放进去这样即便是重复同步Doris也会自动去重。别问我怎么知道的线上报表数据翻倍这种尴尬事经历过一次就长记性了。
返回列表