ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Flink+Hudi实时数仓实战:流式写入、增量读取与数据血缘

Flink+Hudi实时数仓实战:流式写入、增量读取与数据血缘 1. Flink 与 Hudi 为什么要放在一起学先把组合逻辑捋顺我最早接触 Hudi 是在一个离线数仓的改造项目里当时的需求很朴素MySQL 的 binlog 已经通过 Canal 落到 Kafka 了但下游的报表还得等 T1 的 Spark 任务跑完才能看业务方天天催。那会儿试过几套方案最后落在 Flink Hudi 上一直用到现在。所以这篇内容我想按自己的踩坑顺序来讲从为什么选这套组合到建表参数怎么定、增量读怎么配、血缘怎么抽、报错怎么排尽量把能直接抄的部分写清楚。先明确这套组合解决的是什么问题。简单说Flink 负责把源源不断的数据实时搬进来Hudi 负责把搬进来的数据以可更新、可增量读的方式存下来。前者解决的是时效性后者解决的是存储层的可变更能力。单看任何一个都不新鲜难的是拼在一起之后主键索引、小文件、Compaction、Checkpoint 这几件事会互相牵扯配错一个参数作业跑两天就开始堆小文件或者 Checkpoint 超时。适合看这篇的人大致有三类一是刚把 Flink 跑起来、准备接实时数仓的新手想搞清楚 Hudi 到底解决哪一环二是已经在用 Hudi但被 MOR 表的读放大、小文件、Compaction 拖慢折磨的运维同学三是需要把 Flink SQL 作业接入元数据平台、要做数据血缘的架构同学。三类人的关注点完全不同我会在对应章节里分别展开。1.1 离线链路到底卡在哪个环节传统离线数仓的链路是业务库 - 抽数 - Hive 分区表 - Spark 计算 - 结果表。这条链路本身没问题卡住的地方在于分区一旦写完就不愿意改。Hive 表的分区文件是追加写想更新一行历史数据要么重刷整个分区要么在查询时用row_number()取最新版本。前者成本高后者把复杂度全推给了查询层。具体表现就是订单表一天有几十万条状态变更落到 Hive 里就变成几十万条重复记录下游每次查都要去重。数据量小的时候没人管等到单表上万亿行、单分区几十 GB 的时候一次row_number去重能跑四十分钟。这个成本是硬性的跟集群规模没关系因为扫描量摆在那儿。Hudi 的价值就在于把去重这件事从查询时前移到写入时。写入阶段就按主键合并好读的时候直接拿到最新快照或者只读某个时间点之后的增量。这一前一后的差别在多表 join 的场景里会被放大很多倍因为 join 的时候重复数据会被再次放大几亿行的中间结果很容易把 Executor 打爆。1.2 Hudi 补上的三块能力值得单独记一下第一块是主键级别的 upsert。Hudi 的recordkeyprecombine组合决定了同一主键多条记录怎么保留最新这个和数据库的 upsert 语义基本一致但作用在文件层面靠索引定位到具体文件组。第二块是增量读取read.streaming.enabled打开之后Hudi 表可以像 Kafka 一样被 Flink 持续消费只读 commit 时间之后的变更。第三块是时间旅行通过read.start-commit和read.end-commit可以读任意历史快照做回溯、对账、故障恢复都很方便。这三块能力里增量读取是我觉得最容易被低估的。很多团队上了 Hudi 之后还是全量读 Hudi 表那其实只用到了一半价值。真正省资源的是增量读下游的宽表拼接任务每次只消费上游的一小段变更输入量可能只有全量的千分之一Flink 作业的并行度都能降下来。1.3 为什么写入引擎选 Flink 而不是 Spark这个问题我被问过很多次。Spark 也能写 Hudi而且 Structured Streaming 写 Hudi 的文档更成熟为什么还要折腾 Flink核心原因是延迟和状态管理。Spark Structured Streaming 的微批间隔虽然能调到秒级但每个批次的任务调度、JVM 复用、shuffle 开销是省不掉的稳定跑到 10 秒以内的延迟需要花不少力气调优。Flink 是真正的流处理Checkpoint 周期设成 1 分钟端到端延迟就能稳定在分钟级以内而资源占用比微批方式更平滑不会出现批次之间的大起大落。另一个原因是状态和幂等。Flink 的 Checkpoint 机制配合 Hudi 的两阶段提交能做到 Exactly-Once 写入。Hudi 在write.operation为upsert时commit 元数据会记录 checkpoint ID作业重启后从最后一个成功的 checkpoint 恢复不会重复写也不会丢数据。这套机制我在生产上验证过不下十次重启数据没有出现过重复或缺失。还有一个很现实的原因很多团队的实时链路本身就是 Flink 一套Kafka 接入、维表 join、指标计算全是 Flink SQL再单独为写入 Hudi 引一套 Spark 集群运维成本不划算。统一到 Flink 之后作业管理、监控告警、资源调度都能复用。1.4 这套组合适合什么场景不适合什么场景适合的场景我总结成三条判断标准数据有明确主键且会更新、下游对时效要求在分钟到小时级、写入 QPS 不能太高。比如订单状态流水、用户画像标签更新、IoT 设备最新状态、CDC 同步的宽表这些都很合适。QPS 在几千到几万这个量级配合合理的分桶数Hudi 是扛得住的。不适合的场景也说清楚免得踩坑。第一类是只追加不改的日志型数据比如埋点日志、访问日志这种用 Hudi 反而是浪费直接写 Parquet 或者用 Iceberg 更轻因为不需要索引和合并的开销。第二类是超高频点更新比如每秒几十万次的计数器更新Hudi 的文件组写入会被索引查找拖住这种更适合放到 KV 存储里。第三类是超大规模的全量重算Hudi 的 Compaction 和 Clustering 都是重操作如果每天要重刷全部历史成本会很高。提示判断能不能用 Hudi先看主键更新比例。如果更新记录占总记录的 30% 以上Hudi 的合并开销会明显上升这时候要么调大分桶数要么考虑拆表。2. 环境准备版本矩阵、部署模式和那个被问烂的必须 HDFS 吗环境这块踩的坑最多我先说结论Flink 和 Hudi 的版本必须严格配对差一个次版本就可能报 ClassNotFound 或者方法签名不匹配。这个问题不是你代码写得不好纯粹是版本边界问题网上的教程大多不写清楚照着抄很容易翻车。2.1 版本配对是第一道坎先把这张表存下来我这几年维护过的组合大致是这些列出来供对照Hudi 版本支持的 Flink 版本备注0.10.x1.12 / 1.13功能较少不建议新项目用0.11.x1.13 / 1.14引入 Flink SQL Connector 稳定版0.12.x1.14 / 1.15支持 BUCKET 索引写入性能提升明显0.13.x1.15 / 1.16支持流读多表 join 的早期形态0.14.x1.16 / 1.17 / 1.18生产上比较稳的一档1.0.x1.18 / 1.19 / 1.20API 有调整升级需评估选型的建议是新项目直接上 Hudi 1.0 Flink 1.18 以上的组合存量项目如果已经在 0.14 稳定运行不要为了追新去升。我见过一个团队从 0.12 升到 1.0因为表属性不兼容跑了一周的数据全部要重刷代价很大。Java 版本也要注意。Flink 1.18 开始对 Java 17 的支持比较完整但 Hudi 的某些模块在 Java 17 下还有反射相关的告警编译期可能看不出来运行时才暴露。稳妥做法是保持 Java 8 或者 Java 11跟集群上其他组件的 Java 版本一致避免跨版本加载类出问题。2.2 部署模式怎么选Standalone、YARN 还是 K8s本地学习和生产部署的选择完全不同。本地学习用 Standalone 最省事下载 Flink 二进制包start-cluster.sh起来就能提交作业JobManager 和 TaskManager 都在一台机器上交作业的速度很快调试参数的时候改一次跑一次效率高。生产上我推荐YARN Application 模式。理由是资源隔离清楚、Flink 版本互不影响、作业挂了 YARN 会拉起。Session 模式虽然提交快但多个作业共享 JobManager一个作业的 OOM 或者线程泄漏会影响到其他作业线上出过一次因为某作业异常导致整个 Session 集群失联之后我就再没用过 Session 模式跑生产。K8s 模式现在的成熟度也够了优势是资源弹性好、镜像化管理方便特别是你需要频繁调整并行度的时候。缺点是网络和存储的配置复杂度高Hudi 写 HDFS 或者对象存储需要提前把配置挂进镜像调试不如 YARN 直观。团队如果没有成熟的 K8s 运维能力建议先别上。2.3 Hudi Flink Bundle 的编译与依赖坑Hudi 的 Flink Bundle 是需要单独编译的这点很关键。官方发布的 jar 不一定包含你需要的所有依赖尤其是 Hive Sync、对象存储 SDK 这些。编译命令大致是这样mvn clean package -DskipTests -Drat.skiptrue \ -Dflink.version1.18.0 \ -Dscala.binary.version2.12 \ -Dhadoop.version3.3.4 \ -Pflink-bundle-shade-hive3-Pflink-bundle-shade-hive3这个 profile 一定要带上否则 Hive Sync 功能会缺类。编译出来的 jar 在packaging/hudi-flink-bundle/target/目录下放到 Flink 的lib/目录里。依赖冲突是另一个高频问题。常见的是hadoop-common和hadoop-client同时存在导致类加载混乱或者guava版本不一致。排查办法是把 Flink 的lib/目录列一遍同一个 groupId 只保留一个版本特别是hadoop-*、jackson-*、guava这几个。注意如果集群上用 HDFS 做存储Flink 的lib/下必须有匹配的flink-shaded-hadoop-3-uber或者hadoop-client否则提交作业时连不上 NameNode报的是连接超时而不是缺少依赖很容易查错方向。2.4 文件系统抽象对象存储能不能顶替 HDFSFlink 一定要 HDFS 吗这个问题答案是不一定但要看存储层的支持情况。Hudi 底层依赖 Hadoop 的 FileSystem 抽象所以只要是实现了这个接口的存储都能用包括 HDFS、本地文件系统、S3、OSS、COS 这些对象存储。本地学习阶段完全可以用本地路径path写成file:///tmp/hudi/orders跑起来一样能验证功能。生产上选对象存储要注意几点。一是重命名操作HDFS 的 rename 是原子的对象存储的 rename 实际上是 copy delete小文件多的时候 Compaction 会非常慢。二是一致性问题部分对象存储对 list 操作的强一致性支持有限Hudi 的 commit 元数据解析可能读到旧结果。三是元数据操作开销对象存储每次 list 都是网络请求Hudi 的 Clean、Compaction 会频繁 list 文件成本会上去。我的实际选择是实时链路写入用 HDFS归档和冷数据用对象存储。HDFS 的随机读和元数据性能在写入密集场景下优势明显等数据冷下来之后跑一次 Clustering 和归档把老分区挪到对象存储成本能降不少。3. Hudi 表模型选型与 Flink SQL 建表参数详解建表这一步决定了后面所有的性能表现我把它放在最前面讲是因为改表模型比调参麻烦得多选错了要么重建表要么忍受长期的读放大。3.1 COW 和 MOR 怎么选别光看文档COWCopy On Write是写的时候合并每次写入都会重写整个文件组读的时候只有 Parquet 文件读性能好。MORMerge On Read是写的时候只写日志读的时候合并日志和基础文件写快读慢。文档一般会说读多写少用 COW写多读少用 MOR但这个说法太粗。实际选的时候要看读的形态。如果你的下游是 Flink 流读或者 Spark 批读每次只读小部分增量MOR 更合适因为写入延迟低、延迟可见性可以压到分钟级。如果下游是 BI 工具做全表扫描或者大范围聚合COW 更合适因为不用合并日志文件扫描量稳定。我自己的经验是Flink 实时写入 下游也是 Flink 增量读用 MORFlink 实时写入 下游是离线 Spark 全量算用 COW。中间状态比如下游两种都有可以先用 MOR等读放大成为瓶颈再考虑加 Clustering 或者部分表转 COW。3.2 索引类型的选择直接影响写入吞吐Hudi 在 Flink 上支持几种索引最常用的是BUCKET、FLINK_INMEMORY和BLOOM。FLINK_INMEMORY是 Flink 专属的索引把索引维护在 TaskManager 的内存里写入性能最好缺点是内存占用大、作业重启后需要重建索引而且并行度变更后索引失效。适合数据量可控比如单表几亿行以内、并行度稳定的场景。BUCKET索引是在建表时按主键哈希分桶每个桶对应固定的文件组写入时直接定位不需要全局索引查找。它的优势是内存占用低、重启不丢索引、并行度调整友好是目前 Flink Hudi 生产上用得最多的方案。BLOOM索引是基于布隆过滤器的全局索引写入时需要读取所有文件组的布隆过滤器索引开销大一般不推荐在流式写入场景用。选型建议很直接默认用 BUCKET只有在写入吞吐成为明确瓶颈、且并行度长期不变的情况下才考虑 FLINK_INMEMORY。3.3 一份能直接跑的 Flink SQL 建表语句下面这份是我在项目里用了很久的模板参数都是实际调过的可以直接改改路径用CREATE TABLE hudi_order_detail ( order_id BIGINT, user_id BIGINT, product_id BIGINT, order_amount DECIMAL(18, 2), order_status STRING, update_time TIMESTAMP(3), dt STRING, PRIMARY KEY (order_id) NOT ENFORCED ) PARTITIONED BY (dt) WITH ( connector hudi, path hdfs:///warehouse/hudi/order_detail, table.type MERGE_ON_READ, hoodie.datasource.write.recordkey.field order_id, hoodie.datasource.write.precombine.field update_time, hoodie.datasource.write.partitionpath.field dt, write.operation upsert, index.type BUCKET, hoodie.bucket.index.num.buckets 16, hoodie.bucket.index.hash.field order_id, compaction.async.enabled true, compaction.trigger.strategy num_commits, compaction.delta_commits 5, clean.async.enabled true, hoodie.cleaner.commits.retained 10, write.tasks 4, write.batch.size 128, hive_sync.enable false );几个参数值得单独说。precombine.field必须是单调递增的字段一般用update_time或者 binlog 的 ts_ms如果这个字段不单调同一主键的更新顺序会乱数据会错。compaction.trigger.strategy设为num_commits、阈值 5意思是积攒 5 个增量 commit 之后触发一次合并这个值太小会增加合并频率、拖慢写入太大则读放大严重5 到 10 之间是比较舒服的区间。hive_sync.enable我默认关掉。原因是 Hive Sync 会给 JobManager 增加额外的元数据操作作业重启时如果有残留的同步任务会出现锁等待。需要同步的时候单独起一个离线任务做比在线同步稳。3.4 分桶数和写入并行度到底怎么算分桶数是个很实际的问题配少了并发上不去配多了小文件堆成灾。我的计算逻辑是这样先估算单桶的日均数据量控制在 500MB 到 1GB 之间。比如单表日增 20GB 原始数据压缩之后按 1/3 算大概 6~7GB那分 8 到 16 个桶比较合适。这个数字不是死的需要根据写入 QPS 微调如果写入 QPS 很高但单个桶数据量很小可以适当减少桶数让每个桶写得更满。写并行度write.tasks建议跟桶数保持整数倍关系比如 16 个桶配 4 或者 8 个写任务这样每个任务处理的桶数是整数负载均衡。如果桶数是 16、写任务是 5会出现任务之间的数据量差异很大慢任务拖住整个 Checkpoint。提示分桶数一旦设定后续调大需要重写历史数据。所以初期宁可稍微多分一点留出增长空间也不要设得太小。4. 流式写入与增量读取的完整实操参数说完了接下来是实际跑起来的部分。我按数据源准备 - 写入作业 - 增量读 - Compaction 观测的顺序走一遍。4.1 从 Kafka 到 Hudi 的完整链路上游假设是 Kafka 里的订单变更消息格式是 JSON字段跟上面的表结构对应。源表建法CREATE TABLE kafka_order_src ( order_id BIGINT, user_id BIGINT, product_id BIGINT, order_amount DECIMAL(18, 2), order_status STRING, update_time TIMESTAMP(3), proc_time AS PROCTIME() ) WITH ( connector kafka, topic order_change, properties.bootstrap.servers kafka01:9092,kafka02:9092, properties.group.id flink_hudi_order, scan.startup.mode group-offsets, format json, json.ignore-parse-errors true );scan.startup.mode用group-offsets配合 Checkpoint 做精确一次消费。如果作业第一次启动会从 group 的 committed offset 开始做数据回灌的时候才改成earliest-offset。json.ignore-parse-errors建议打开。线上 JSON 格式偶尔会有脏数据比如字段类型不对或者嵌套层级异常不打开的话整个作业会因为一条脏数据挂掉打开了至少能跳过这条继续跑同时在日志里留下告警。写入语句就是把源表的数据打进 Hudi 表注意分区字段需要从事件时间派生INSERT INTO hudi_order_detail SELECT order_id, user_id, product_id, order_amount, order_status, update_time, DATE_FORMAT(update_time, yyyy-MM-dd) AS dt FROM kafka_order_src;这里的dt用的是事件时间不是处理时间。用处理时间会导致迟到数据落到错误的分区跨天的时候特别明显。4.2 Checkpoint 配置对写入稳定性的影响Checkpoint 周期直接决定了可见延迟和数据一致性窗口。我一般设60 秒理由是周期太短会增加 HDFS 的元数据压力和 Hudi 的 commit 频率周期太长则作业恢复时要重放的数据多而且 Hudi 的 commit 间隔变大下游增量读的延迟跟着涨。配置写在flink-conf.yaml或者提交参数里execution.checkpointing.interval: 60s execution.checkpointing.timeout: 10min execution.checkpointing.min-pause: 30s execution.checkpointing.max-concurrent-checkpoints: 1 execution.checkpointing.externalized-checkpoint-retention: RETAIN_ON_CANCELLATION state.backend: rocksdb state.backend.incremental: truemin-pause设成周期的一半避免两个 Checkpoint 背靠背跑造成资源争抢。max-concurrent-checkpoints一定要设成 1同时跑多个 Checkpoint 在写 Hudi 的场景下容易出现文件写冲突反而更慢。state.backend.incremental打开之后RocksDB 的 Checkpoint 只上传增量部分状态大的作业恢复速度能快好几倍。4.3 增量读取把全量扫描变成小段消费Hudi 表建好之后可以直接建一个流读表CREATE TABLE hudi_order_inc ( order_id BIGINT, user_id BIGINT, order_status STRING, update_time TIMESTAMP(3), dt STRING ) WITH ( connector hudi, path hdfs:///warehouse/hudi/order_detail, table.type MERGE_ON_READ, read.streaming.enabled true, read.start-commit earliest, read.streaming.check-interval 1, read.streaming.skip_compaction true, read.tasks 2 );read.streaming.check-interval是轮询新 commit 的间隔单位分钟设成 1 表示每分钟检查一次。read.streaming.skip_compaction打开后读的时候只消费增量 commit不消费 Compaction 产生的 commit避免同一条数据被读两次——这个参数非常关键不打开的话会出现数据重复。增量读的实际收益很直观。之前一个宽表拼接的任务全量读每天要扫 300GB 数据改成增量读之后每次只消费几 MB 到几十 MB 的变化Flink 的并行度从 16 降到 4资源成本直接砍掉三分之一。4.4 Compaction 的调度和观测方式MOR 表如果 Compaction 不跑日志文件会无限增长读的时候要合并的日志越来越多最后读性能崩掉。调度方式有两种在线异步 Compaction 和离线定时 Compaction。在线方式就是建表时的compaction.async.enabled true由 Flink 作业自己在写入过程中触发。好处是自动化程度高缺点是抢占写入资源高峰期会拖慢写入延迟。离线方式是单独起一个 Flink 批作业做 Compaction时间可以选在业务低峰期对主链路没影响。观测 Compaction 状态可以看 Hudi 的 commit 时间线用命令行工具hudi-cli connect --path hdfs:///warehouse/hudi/order_detail commits show --sortBy Total Time Taken compactions show all重点看两个指标未合并的增量 commit 数量和单个文件组的大小。前者超过阈值就说明 Compaction 落后了后者超过 1GB 说明需要 Clustering 来重新组织文件布局。注意Compaction 失败不会自动重试需要手动重启。如果发现时间线上有 pending 状态的 Compaction先查清楚失败原因再重跑直接重跑可能会因为文件状态不一致再次失败。5. 数据血缘从 Flink SQL 里把表级和字段级关系抽出来血缘这件事在离线数仓里相对简单因为 SQL 是静态的解析一次就行。到了 Flink SQL作业是长期运行的而且流式作业的表关系会随着 DDL 变更而变化所以需要一套持续采集的机制。5.1 为什么实时链路的血缘更值得做实时链路的数据流向比离线长一条数据可能从 Kafka 进 Flink写进 Hudi再被下游 Flink 增量读出来写进另一张 Hudi 表或者 Doris中间还可能经过多次维表 join。出了问题要定位是哪个环节的数据错了如果没有血缘只能一个个作业去翻 SQL效率极低。还有一个更现实的场景是影响面分析。上游表结构调整一列或者某个字段的语义变了需要知道会影响多少下游作业。离线场景下可以等调度系统报错再说实时场景下错误数据会持续流入可能几个小时才被发现这时候血缘就是止损的关键工具。5.2 基于 Calcite 解析的字段级血缘思路Flink SQL 底层用的是 Calcite 做解析和优化所以血缘可以沿着这条链路来取。第一个层次是表级血缘相对简单。通过 Calcite 解析 SQL 得到 AST遍历SqlInsert、SqlSelect节点提取源表和目标表。对于 INSERT INTO SELECT 这种结构源表集合来自 FROM 和 JOIN目标表来自 INSERT 的 target。这套解析逻辑不复杂自己写一个 Visitor 就能拿到。第二个层次是字段级血缘需要用到 RelNode。Calcite 在 validate 之后会生成 RelNode 树叶子节点是TableScan中间是Project、Filter、Join这些。字段的映射关系藏在Project节点的RexNode里遍历表达式就能建立 目标字段 - 源字段 的映射。碰到SELECT *或者UNION这类的场景需要做展开和合并稍微麻烦一点但逻辑是通的。第三个层次是跨作业串联。单个作业的血缘只解决了一半问题真正有用的是把多个作业的血缘拼起来形成一条从 Kafka 到最终结果表的完整链路。这就需要给每个作业打上唯一标识把解析结果写进元数据平台然后按表名做关联。5.3 落地过程中的几个实际约束第一DDL 要先采集。SQL 解析依赖表的 schema 信息如果解析时拿不到源表的字段定义字段级血缘就建不起来。所以需要先通过 Flink 的 Catalog API 把源表、目标表的 DDL 采集下来存进元数据库。第二函数和 UDF 的处理。SQL 里用了自定义函数的话Calcite 解析出来只能看到函数名看不到内部逻辑。这种血缘要标注成不透明不要硬猜。第三增量采集。实时作业的 SQL 可能因为业务调整而改动血缘需要有版本概念记录每次变更的时间和内容否则追溯历史问题时对不上。我自己做的方案是作业提交时触发一次 SQL 解析结果写入元数据表表结构是(job_id, source_table, target_table, column_mapping, version, create_time)。血缘查询就是在这张表上做递归查询找到所有下游节点。这个方案不复杂但确实解决了很多排查问题的时间。6. 高频异常排查速查JDBC、Doris、TiDB CDC 都在这这一节是踩坑记录我把实际遇到过的几类问题整理出来附上排查思路。这些问题在搜索引擎上被问得很多但答案往往零散我按现象 - 原因 - 解决整理成表方便对号入座。报错现象可能原因处理方式No suitable driver found驱动 jar 未放入 lib 或 classpath确认驱动 jar 在lib/下且版本与目标库匹配Connection is not available, request timed out连接池被打满或连接泄漏调大连接池上限检查是否有流未关闭Packet for query is too large批量写入单条过大调小sink.buffer-flush.max-rows或调大服务端max_allowed_packetdatev2 but arrow type is dateday类型映射表缺项升级连接器版本或统一时间类型定义6.1 Flink JDBC 连接器异常的排查顺序JDBC 相关问题我一般按五步查。第一步看驱动lib/下有没有对应的 JDBC 驱动 jar版本是否匹配这块能解决一半以上的问题。第二步看连接字符串特别是时区和字符集参数serverTimezoneAsia/Shanghai这类参数缺失会导致时间字段差 8 小时表象是数据对不上但作业不报错很难发现。第三步看连接池。Flink JDBC Sink 默认用的是连接池如果写入速率高于连接获取速率池子会被占满。可以调大connection.max-retry-timeout和池子大小但更根本的办法是看写入速率是不是超过了目标库的承受能力。第四步看批量参数。sink.buffer-flush.max-rows和sink.buffer-flush.interval这两个决定了攒批策略。设得太激进比如 10000 行一批容易触发包大小限制设得太保守则写入效率低。我一般设max-rows为 500 到 1000interval为 1 到 2 秒平衡吞吐和延迟。第五步看目标库侧。锁等待、慢查询、连接数上限这些在 Flink 侧看不出来需要到数据库看监控。有一次排查了很久最后发现是目标库上有个长事务持有表锁导致写入全部阻塞。提示JDBC Sink 如果配了主键但目标表没有对应约束会出现主键冲突报错。检查两边的主键定义是否一致这个不一致在测试环境很难暴露上线后才出现。6.2 Doris Connector 的 datev2 类型映射报错flink type is datev2, but arrow type is dateday这个报错是因为 Doris 在较新版本里把 DATE 类型的内部表示换成了 DATEV2而 Flink Connector 的类型映射表里还是按老的 DATEDAY 来处理匹配不上就走了兜底分支直接抛异常。处理方式有三个方向。升级连接器版本是最推荐的新版本的flink-doris-connector已经修了这个映射问题选跟 Doris 服务端版本对应的连接器即可。降级 Doris 的列类型是临时方案把相关列的建表语句从 DATE 改成 DATEV1 的写法但这个在生产上需要改表代价不小。显式指定 schema也算一个规避办法在 Flink 建表时不依赖自动推导手动把字段类型写成 Flink 的 DATE 或者 TIMESTAMP跳过 Doris 侧的类型识别。这个报错常见于用table.identifier读取 Doris 表结构的场景因为这时候 schema 是从 Doris 的元数据里拉的会带上 DATEV2 的信息。如果你只是往 Doris 写数据不走自动 schema 推导一般不会碰到。6.3 TiDB CDC 接入 Flink SQL 的注意事项从 TiDB 接 CDC 到 Flink链路一般是 TiCDC 采集 - Kafka - Flink 消费或者直接用flink-connector-tidb-cdc直连。两种方式各有取舍。走 Kafka 的方式链路长但解耦性好TiCDC 的稳定性不直接影响 Flink 作业出问题可以重放 Kafka。直连方式延迟低但需要 Flink 作业持续持有 TiKV 的连接作业重启期间如果 TiCDC 的 GC 水位推进了会读不到变更数据。实际配置里要注意几个点。GC 水位TiDB 的 GC 会清理历史版本如果 Flink 作业因为故障停了很久恢复时的 checkpoint 时间点早于 GC 水位就会报 changelog 丢失。解决办法是适当调大 GC 保留时间同时给作业加恢复告警不要让它停太久。scan.startup.mode的选择首次启动用initial做全量增量之后用latest-offset只消费增量不要每次都用initial否则每次重启都全量扫一遍。表结构变更的处理也要考虑。TiDB 上加列之后Kafka 里的消息结构会变Flink 的源表如果没跟着改解析会失败。稳妥做法是让 Flink 的源表字段是目标结构的超集多出来的字段留空这样加列不用重启作业。6.4 小文件、OOM 和 Checkpoint 超时的连锁反应这三个问题经常一起出现本质上是同一件事的不同表现。小文件的根源是写入并行度高但单次写入数据量小。每个 Checkpoint 周期都会为每个桶生成一批新文件如果周期短、数据少文件就会碎。缓解办法是适当调大 Checkpoint 周期、增大write.batch.size、降低写并行度。已经产生的小文件用 Clustering 来整理配置离线任务定期跑。OOM一般出在 TaskManager 上原因是状态太大或者索引占用内存高。用FLINK_INMEMORY索引的话内存占用会随数据量线性增长这也是我推荐BUCKET的原因。如果是状态太大检查一下是否 join 了大表没设 TTL状态不过期就会一直涨。Checkpoint 超时往往是前两个问题的结果。小文件多了之后Hudi 的 commit 需要处理的文件数量增加超时时间被拉长内存压力大了之后 GC 频繁Checkpoint 的 barrier 对齐变慢。所以遇到 Checkpoint 超时不要只调大超时时间要往上游看是不是小文件或者状态的问题。指标健康区间异常时的动作单文件组大小100MB ~ 1GB超过则跑 Clustering未合并 commit 数小于 10超过则检查 Compaction 任务Checkpoint 耗时小于周期的 50%超过则排查状态和文件数写任务反压率小于 20%超过则调整并行度或桶数7. 关于参数调优我自己的几条判断标准调参这件事文档给的是起点真正的数字得在自己的数据和集群上试出来。我个人的习惯是先固定变量一次只调一个参数跑满一天看指标不要一次改五个参数出了问题不知道是哪个的影响。有几个经验值可以当参考。Checkpoint 周期从 60 秒起调如果发现 commit 次数太多导致元数据压力大可以放到 120 秒如果可见性要求高往下调到 30 秒但要注意小文件。write.batch.size从 128 起调写入吞吐上不去可以往 256 加但要盯着目标存储的随机写能力HDFS 对这种批量写是友好的对象存储就未必。Compaction 的delta_commits从 5 起调读放大明显就降到 3写入延迟紧张就升到 10。还有一个容易被忽略的点是时间字段的精度。TIMESTAMP(3)和TIMESTAMP(6)在 Hudi 里的存储和比较行为不一样如果上游的update_time是毫秒精度建表用了TIMESTAMP(6)precombine 的比较结果可能会有细微差别。统一用TIMESTAMP(3)是最省事的做法除非业务确实需要微秒级。最后分享一个排查思路遇到数据不一致先看时间线再看文件最后看 SQL。Hudi 的时间线记录了每一次 commit 的状态和时间如果某个 commit 的状态是 pending 或者 failed问题基本就定位到了。时间线没问题再去看对应分区的文件列表看有没有文件大小异常或者数量异常。这两步排完还没找到原因再去检查 SQL 逻辑和参数配置顺序别反了不然会浪费很多时间。
返回列表