全文围绕 Flink 与 Hudi 的学习路径展开,从环境搭建到 Flink SQL、Hudi 表模型、下游对接与血缘治理,逐步拆解,力求给出一份可抄作业的实操参考。
1. 我为什么把 Flink 和 Hudi 绑在一起学
刚接触实时数仓那会儿,我的认知是割裂的:Flink 是算的,Hudi 是存的,两个东西各学各的,好像没什么交集。直到真正上手做一条从业务库到分析层的实时链路,才发现这两个东西根本拆不开——Flink 负责把数据算清楚,Hudi 负责把结果存明白,中间任何一环掉链子,整条链路都是半成品。所以我把 Flink & Hudi 当成一个整体来学,不是两个独立的知识点,而是一套组合拳。
这篇文章面向的是已经会写点 SQL、对大数据组件有初步印象,但一到“实时链路怎么落地”就发懵的同学。我会把环境怎么搭、Flink SQL 怎么写、Hudi 表怎么建、参数为什么这么调、踩过哪些坑,尽可能摊开讲清楚。有基础的可以直接跳到对应章节抄配置,新手建议按顺序读一遍,因为后面很多参数的取舍,逻辑都在前面几节里。
1.1 离线数仓留下的断层到底在哪
传统离线链路的结构大家都很熟:业务库通过抽取工具同步到分布式文件系统,再按天分区落到 Hive 表里,跑完调度任务,第二天早上出报表。这套东西稳定、成本低、生态成熟,但它有一个绕不过去的硬伤——时效性被锁死在 T+1。
问题在于,业务侧的需求早就不满足于“看昨天的数”了。风控要秒级识别异常交易,运营要实时看大促期间的成交曲线,推荐系统要拿分钟级的用户行为特征。你让这些场景等一天,等于没做。于是大家开始往实时方向找方案,第一反应通常是“那我把抽取频率调高不就行了”。实测下来,这条路走不通:每小时全量抽一次,数据库直接被拖垮;改成增量抽,又面临更新和删除怎么表达的问题——离线表是按分区覆盖写的,一条数据改了,你没法只改那一行。
这就是断层所在。离线存储模型的写入语义是“批量覆盖”,而实时场景需要的是“行级变更”。Flink 解决了计算侧的实时性,但如果没有一个能承接行级更新、又能被下游高效查询的存储层,算得再快也没地方落。
1.2 Hudi 在链路里补的是哪块板
Hudi 的价值,就是把这最后一块板补上。它在分布式文件系统之上提供了一层表抽象,核心能力有三点:支持按主键做 upsert、支持增量读取、支持对文件做自动治理。
打个比方,把分布式文件系统想象成一个巨大的仓库,Hudi 就是仓库里的智能货架系统。原始文件系统只会告诉你“这堆货是什么时候放进去的”,而 Hudi 会记录“哪个货位上的哪件货被谁在什么时候换过”,并且定期把零散的小包裹合并成大托盘。对 Flink 来说,写入 Hudi 就像往一个支持主键更新的数据库里写数据;对下游的查询引擎来说,读 Hudi 表跟读普通表没什么区别。
我最看重的是它的增量读能力。以前做实时链路,ODS 层到 DWS 层要么全量重算,要么靠消息队列再串一次,链路又长又脆。有了 Hudi,DWS 层可以直接读 ODS 层过去 5 分钟的增量,计算量只有全量的百分之几,这一下就把资源成本压下来了。
1.3 这套组合适合谁,不适合谁
先说适合的。数据量在千万到百亿级别、有明确主键、需要分钟级甚至秒级时效、同时下游既要明细查询又要聚合分析的场景,Flink + Hudi 是很对路的。电商订单、物流轨迹、用户行为埋点、IoT 设备上报,这些都属于典型适用场景。
再说不太适合的。如果你的数据压根没有稳定主键,比如纯日志类文本,那 Hudi 的 upsert 能力用不上,直接写对象存储或者用普通分区表更省事。如果数据量只有几十万行,用传统数据库加个定时任务就够了,上这套组合属于拿高射炮打蚊子。还有一种情况是团队里没人懂分布式文件系统和集群运维,那前期投入的学习成本会很高,建议先用托管服务过渡。
实操心得:判断要不要上 Hudi,我一般只问两个问题——第一,数据是否有明确且稳定的主键;第二,下游是否存在“既要实时又要能改历史”的需求。两个都是“是”,就值得上;有一个是“否”,先掂量掂量。
2. 环境搭建:从 Flink 安装配置到部署这条线怎么走顺
环境这一节我放在最前面讲,是因为太多人卡在这里就放弃了。Flink 的安装本身不复杂,复杂的是版本匹配和依赖关系。下面按“先定版本、再搭单机、然后选部署形态、最后准备存储”这个顺序来,每一步我都会说清楚为什么这么做。
2.1 版本矩阵先定死,别边装边试
这是我最想强调的一条经验:动手之前先把版本矩阵写在一张纸上,后面所有操作都围绕它来。Flink、Hudi、分布式文件系统、Java 版本之间是有兼容关系的,尤其是 Hudi 和 Flink 的绑定版本,错一个小版本就可能报一堆看不懂的类找不到错误。
下面这张表是我在几个项目里验证过相对稳妥的组合,供参考:
| 组件 | 版本 | 说明 |
|---|---|---|
| JDK | 8 或 11 | Hudi 对 17 支持一般,建议 11 |
| Flink | 1.14 / 1.16 / 1.17 | 1.16 之后 Hudi 集成更顺 |
| Hudi | 0.12 / 0.13 / 0.14 | 必须严格对应 Flink 小版本 |
| Hadoop | 3.2 / 3.3 | 只作为存储客户端,不一定要自己装 |
| Scala | 2.12 | 与 Flink 编译版本一致 |
为什么强调小版本严格对应?因为 Hudi 的 Flink 集成包是hudi-flink1.14-bundle、hudi-flink1.16-bundle这种命名方式,它是按 Flink 大版本单独编译的,内部依赖的 Flink API 签名不一样。你拿 1.14 的包丢进 1.16 的环境里,编译期可能不报错,运行期直接给你抛NoSuchMethodError,排查起来非常痛苦。
注意:下载依赖包时,认准 bundle 包。Hudi 社区同时提供
hudi-flink-bundle和一堆独立 jar,初学者直接上 bundle,别拆开一个个凑,省下来的时间够你多跑好几个作业。
2.2 单机 Standalone 最小可用集群
正式集群之前,我强烈建议先用一台机器跑通 Standalone 模式。目的不是性能,是让你把配置文件里的每一项都摸一遍。
下载解压之后,主要改两个文件。第一个是conf/flink-conf.yaml:
jobmanager.rpc.address: localhost jobmanager.memory.process.size: 2048m taskmanager.memory.process.size: 4096m taskmanager.numberOfTaskSlots: 4 parallelism.default: 4 # 检查点配置,学习阶段先用文件系统 execution.checkpointing.interval: 60s execution.checkpointing.mode: EXACTLY_ONCE execution.checkpointing.timeout: 10min state.backend: rocksdb state.backend.incremental: true state.checkpoints.dir: file:///data/flink/checkpoints state.savepoints.dir: file:///data/flink/savepoints # 历史服务器,看作业执行图必备 historyserver.web.address: 0.0.0.0 historyserver.web.port: 8082第二个是conf/workers,单机就填一个 localhost。然后./bin/start-cluster.sh启动,浏览器打开 8081 端口能看到界面,说明成了。
这里的参数不是随便填的。taskmanager.numberOfTaskSlots表示一个 TaskManager 能跑几个并行子任务,它和parallelism.default配合决定资源分配。我见过有人把 slot 数设成 1,然后作业并行度设 8,结果八个子任务挤在一个槽里排队跑,还以为是数据倾斜。简单记:一个 slot 大约吃 1 到 1.5 GB 内存,根据 TaskManager 总内存倒推 slot 数,4 GB 的 TaskManager 给 4 个 slot 是比较安全的。
state.backend选 rocksdb 是因为它把状态存在本地磁盘而不是堆内存里,状态大的作业不容易把 JVM 撑爆。配合state.backend.incremental: true,每次检查点只上传变化的部分,检查点时间能缩短一大截。
2.3 部署形态选型:Standalone、YARN、Kubernetes 怎么挑
学习阶段用 Standalone 够了,但真上生产得选形态。我把三种常见形态的取舍整理成下面这张表:
| 形态 | 资源调度 | 适用场景 | 主要痛点 |
|---|---|---|---|
| Standalone | 自管 | 小规模、固定作业、测试 | 资源无法共享,作业多了要手工扩 |
| YARN | 统一调度 | 已有 Hadoop 体系,批流混跑 | 资源竞争时作业可能被抢占 |
| Kubernetes | 容器编排 | 云原生环境,弹性要求高 | 需要额外的 Operator 与镜像管理 |
选型逻辑其实很简单:看你现有的基础设施。如果公司已经有 Hadoop 集群并且在上面跑 Spark 任务,那就上 YARN,不用重复建设;如果基础设施是容器化的,直接上 Kubernetes,配合 Flink Operator 管理作业生命周期;如果只是两三个固定作业、数据量也不大,Standalone 反而最省心,别为了“先进”给自己找麻烦。
我个人的偏好是 Kubernetes。原因不在于技术多先进,而在于作业的部署和回滚变成了标准化的镜像操作,配合 CI 流程可以做到改一行配置就自动发布,这对多人协作的团队很关键。不过前提是你得先有镜像仓库、有资源配额管理,这些前置条件缺一个都别急着上。
2.4 存储层:Flink 落地 Hudi 为什么绕不开分布式文件系统
这个问题被问得特别多:Hudi 能不能不依赖分布式文件系统,直接写本地盘或者对象存储?
技术上,本地盘是能跑的,学习测试完全可以。但生产环境基本不行,原因有两个。一是本地盘容量有限且没有副本机制,一台机器挂了数据就没了;二是 Flink 是分布式计算,多个 TaskManager 跑在不同机器上,写本地盘意味着数据散落在各个节点,下游读的时候得把机器列表都拼起来,根本不现实。
所以 Hudi 需要一个所有计算节点都能访问、具备副本容错、支持大文件顺序写的共享存储,这正是分布式文件系统擅长的事。对象存储同样满足这个条件,而且成本更低,只是它的重命名操作代价较高,会影响 Hudi 的提交延迟,需要用一些参数去适配。
如果你只是想先跑通,可以用file:///的本地路径,但一定要知道这只是学习用。等要把作业带到集群上跑,就得把路径换成hdfs://或者对象存储的地址,同时把依赖的存储客户端 jar 放进 Flink 的lib目录,否则会报No FileSystem for scheme这类错误——这个错误几乎是所有人第一次上集群必踩的坑。
3. Flink SQL 上手:把实时链路写成 SQL 的性价比
学会 DataStream API 之后,我一度觉得写 SQL 是“偷懒”。后来做业务迭代,需求三天两头改一次字段,用 DataStream 每改一次都要重新编译打包上线,而 Flink SQL 改一段建表语句加一句 INSERT 就完事了。从那以后,能在 SQL 层解决的问题我绝不下沉到代码。这一节就把 Flink SQL 的实操骨架讲透。
3.1 一个作业的最小骨架
Flink SQL 的核心就三件事:声明源表、声明目标表、写 INSERT 语句。中间的所有转换逻辑,都在 INSERT 里用 SELECT 表达。
源表从消息队列读订单变更:
CREATE TABLE kafka_order ( order_id STRING, user_id BIGINT, amount DECIMAL(18,2), order_status STRING, update_time TIMESTAMP(3), dt STRING, WATERMARK FOR update_time AS update_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'order_binlog', 'properties.bootstrap.servers' = 'kafka01:9092,kafka02:9092', 'properties.group.id' = 'g_order_hudi', 'scan.startup.mode' = 'group-offsets', 'format' = 'json', 'json.ignore-parse-errors' = 'true', 'json.fail-on-missing-field' = 'false' );目标表落到 Hudi:
CREATE TABLE ods_order ( order_id STRING, user_id BIGINT, 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/ods/ods_order', 'table.type' = 'MERGE_ON_READ', 'write.operation' = 'upsert', 'hoodie.datasource.write.precombine.field' = 'update_time', 'index.type' = 'BUCKET', 'hoodie.bucket.index.num.buckets' = '8', 'write.tasks' = '4', 'compaction.async.enabled' = 'true', 'compaction.schedule.enabled' = 'true', 'compaction.tasks' = '2' );写入就是一句:
INSERT INTO ods_order SELECT order_id, user_id, amount, order_status, update_time, dt FROM kafka_order;注意PRIMARY KEY ... NOT ENFORCED这个写法。NOT ENFORCED表示 Flink 不做主键唯一性校验,只是告诉连接器“这个字段是主键,请你按主键去处理”。这个是必须写的,不写的话 Hudi 连接器不知道拿哪个字段做去重主键,只能退化成插入模式,重复数据会越积越多。
3.2 时间语义与水位线,别等到数据错乱才补
水位线这部分,很多人是先照着模板写、出问题了才回头理解。我建议一开始就搞明白,因为它直接决定窗口计算什么时候触发。
WATERMARK FOR update_time AS update_time - INTERVAL '5' SECOND这句话的含义是:系统认为当前收到的事件时间比真实时间晚最多 5 秒,等水位线推进到某个时间点,就认为该时间点之前的数据都到齐了。这个 5 秒就是允许的最大乱序时间。
设大了,窗口结果出得慢;设小了,迟到的数据会被丢掉。怎么定?看你的数据源乱序程度。同机房消息队列同步过来的数据,1 到 2 秒够了;跨地域汇聚的数据,可能要 10 秒以上。我一般先设保守值,跑一周看迟到数据的分布,再往下调。
还有一个容易忽略的点:如果只有处理时间没有事件时间,就不需要水位线。处理时间写起来简单,但重跑作业的结果不可复现;事件时间需要水位线,但结果稳定可回溯。生产环境的作业,我基本都用事件时间。
3.3 JDBC 连接器异常排查实录
flink的jdbc连接器异常是搜索量很高的词,我自己也踩过不少。这类异常的表象五花八门,但根因基本集中在四类,我按遇到频率排一下。
第一类是驱动缺失。报错通常长这样:java.lang.ClassNotFoundException: com.mysql.cj.jdbc.Driver。原因是 Flink 的lib目录里没有对应数据库的 JDBC 驱动 jar。解决办法很简单,把驱动包丢进lib,重启集群。要注意的是驱动版本要和数据库服务端匹配,用 5.x 的驱动连 8.x 的库,会报 SSL 或者认证插件相关的错。
第二类是连接数打满。报Too many connections,是因为 Flink 的并行子任务每个都会建自己的连接池。假设并行度 8,每个子任务池大小 5,那就是 40 个连接起步,再加上多个作业,很容易把数据库的max_connections撑爆。我的做法是把sink.buffer-flush.max-rows调大、sink.buffer-flush.interval适当延长,用批量换连接数,同时给连接池加空闲回收。
CREATE TABLE dws_user_amount ( user_id BIGINT, dt STRING, total_amount DECIMAL(18,2), PRIMARY KEY (user_id, dt) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:mysql://db01:3306/dw?useSSL=false&rewriteBatchedStatements=true', 'table-name' = 'dws_user_amount', 'username' = 'dw_writer', 'password' = '${secret_value}', 'sink.buffer-flush.max-rows' = '2000', 'sink.buffer-flush.interval' = '3s', 'sink.max-retries' = '3', 'connection.max-retry-timeout' = '60s' );第三类是主键冲突导致写入失败。JDBC 连接器执行的是 UPSERT 语义,底层是 REPLACE 或 ON DUPLICATE KEY UPDATE,如果目标表的主键定义和 Flink 侧声明的不一致,就会出现“写进去但数据不对”或者直接报错。一定要核对两边主键字段顺序和类型完全一致。
第四类是事务超时。批量写大事务时,数据库端wait_timeout或innodb_lock_wait_timeout到了就把连接掐了,Flink 侧报连接已关闭。办法是缩短批量间隔、减小批量行数,让事务小而快。
注意:
${secret_value}这种写法是占位符,真实环境中千万不要把密码明文写进 SQL 文件再提交到版本库。用 Flink 的配置项注入,或者接密钥管理服务。
3.4 跨引擎类型映射:datev2 和 dateday 对不上的时候
有一类错误特别有迷惑性,比如在把 Flink SQL 的结果写进分析型数据库时报flink type is datev2, but arrow type is dateday。第一次看到这个我愣了半天,两个类型名看着都像日期,怎么就冲突了。
本质上这是两个系统对日期类型的不同表达方式发生了碰撞。Flink SQL 侧声明的字段类型经过连接器转换后,变成了一种带精度语义的日期类型,而下游数据库实际列的类型是另一种更宽泛的日期表达,双方的“字典”对不上。中间还有一个传输层用了列式格式做数据交换,它自己又有第三套类型定义,三方各说各话,就报错了。
解决思路有三条,按优先级来:
第一,让两边类型对齐。把 Flink SQL 建表语句里该字段的类型,改成和下游实际列类型完全一致的那一种。如果下游是DATE,那 Flink 侧就用DATE;如果下游是带时分秒的,就用TIMESTAMP,别用TIMESTAMP(3)这种带精度声明的去试探。
第二,在中间加一层转换。用CAST显式转换,把类型语义固定下来:
INSERT INTO doris_daily_report SELECT dt, CAST(stat_date AS DATE) AS stat_date, CAST(amount AS DECIMAL(20,2)) AS amount FROM dwd_daily_agg;第三,降级传输格式。列式传输对类型要求严格,如果实在对不齐,改成按行传输的文本格式,容错空间会大很多。代价是性能下降一些,但对于字段不多、数据量中等的表,这点损失完全可以接受。
实操心得:跨引擎的类型问题,我的排查顺序永远是“先看两边建表语句的类型声明是否完全一致,再看传输格式,最后才动手改代码”。十次里有八次问题就出在第一步,改建表语句比改代码快得多。
4. Hudi 表模型与流式写入实操
环境通了、Flink SQL 会写了,就该把数据真正落到 Hudi 上。这一节讲表模型怎么选、参数怎么读、写入链路怎么搭,以及后期数据治理怎么做。
4.1 COW 和 MOR,先想清楚读放大与写放大
Hudi 有两种表类型:写时复制(COW)和读时合并(MOR)。这两个名字很直白,但选错了会很痛苦。
COW 的逻辑是:每次有更新,就把整个数据文件重写一遍,生成新的文件版本。好处是读的时候非常快,因为它只有一种文件形态,直接读就行。代价是写放大严重,改一行数据可能要重写 128 MB 的文件。
MOR 的逻辑是:更新先写进一个增量日志文件,读的时候再把日志和数据文件合并。好处是写很快,适合高频更新。代价是读的时候要做合并,延迟高一些,而且必须定期做压缩(compaction),否则日志文件越堆越多,读性能会崩。
怎么选?我的判断标准很简单:看读写比。读多写少、并且要求查询延迟低的场景,选 COW,比如面向报表的明细表;写多读少、或者更新极其频繁的场景,选 MOR,比如从业务库同步过来的原始表,这类表更新频繁但很少被直接查询。
有一个折中方案很多人不知道:MOR 表可以配置成近实时读,也就是查询时只读已经压缩过的部分,忽略未压缩的日志。这样既保住了写入速度,又让一部分查询走快路径。代价是数据有分钟级延迟,适合对时效不敏感但对性能敏感的场景。
4.2 建表参数逐个拆解
Hudi 的参数确实多,但真正影响性能和正确性的就那么几个。我把最关键的几项列出来,并解释为什么要这么设:
| 参数 | 取值示例 | 作用与取舍 |
|---|---|---|
table.type | COW / MOR | 表类型,决定读写放大走向 |
write.operation | upsert / insert / bulk_insert | 写入语义,需要主键时用 upsert |
precombine.field | update_time | 同主键多版本时取最大的那条 |
index.type | BUCKET / FLINK | 索引类型,BUCKET 适合大表 |
bucket.index.num.buckets | 8 | 桶数,决定写入并行度和文件数 |
write.tasks | 4 | 写入并行任务数 |
compaction.tasks | 2 | 压缩并行任务数 |
compaction.async.enabled | true | 异步压缩,避免阻塞写入 |
precombine.field这个参数值得单独说。它解决的问题是:同一条订单在两秒内被改了两次,消息队列里是两条消息,都可能被消费到。Hudi 在这两条记录里比较update_time,只保留时间较大的那条。如果这个字段选错,比如选了业务时间而不是更新时间,就会出现旧数据覆盖新数据的情况,而且这种错误很隐蔽,往往要等对账时才发现。
bucket.index.num.buckets是另一个关键。桶的本质是把数据按主键哈希分到若干个固定的文件组里,桶数一旦确定就不能改了(改了就相当于重新分桶)。桶数太少,单个桶数据量大,写入并行度上不去;桶数太多,小文件满天飞。我的经验公式是:目标单文件大小按 128 MB 算,桶数 = 预估数据总量 / 单桶容量,然后向上取整到 2 的幂次,方便后续调整。比如一天 10 GB 数据、保留 30 天、目标单桶 256 MB,那桶数大概在 1200 左右,实际可以取 1024 或 2048。
4.3 Flink 流式写入 Hudi 的完整链路
完整的链路是这样的:消息队列作为源头,Flink SQL 消费并做轻度清洗,写入 Hudi 的 ODS 层;再用一个作业读 ODS 的增量,做聚合后写进 DWS 层。两个作业串联,各管一段。
DWS 层的聚合作业,关键是用增量读而不是全量读:
CREATE TABLE ods_order_incremental ( order_id STRING, user_id BIGINT, amount DECIMAL(18,2), update_time TIMESTAMP(3), dt STRING ) WITH ( 'connector' = 'hudi', 'path' = 'hdfs:///warehouse/ods/ods_order', 'table.type' = 'MERGE_ON_READ', 'read.streaming.enabled' = 'true', 'read.streaming.check-interval' = '60', 'read.start-commit' = 'earliest' );read.streaming.enabled打开流式读,read.streaming.check-interval决定多久检查一次新提交,单位是秒。设成 60 就是每分钟拉一次增量。这个值不要设太小,比如设成 1 秒,会导致大量空轮询,把存储的元数据服务压垮。
写入和读取是解耦的,两边各按自己的节奏来。这个特性对稳定性帮助很大:如果 DWS 层作业挂了要修,修好之后从上次的提交点继续读就行,不会丢数据。
启动这两个作业的顺序也有讲究。先起 ODS 写入作业,跑几分钟确认有数据落盘;再起 DWS 增量读作业。反过来做的话,DWS 作业会一直等不到数据,日志里全是空轮询,看着像出问题了。
4.4 小文件、Compaction 与 Clustering 的治理节奏
Hudi 用久了,绕不开文件治理。这块我刚开始完全没在意,直到某次查询突然变慢,去看目录发现一个分区里有几千个几百 KB 的文件,元数据加载直接卡住。
Compaction 解决的是 MOR 表的日志堆积。开启异步压缩后,Flink 作业会在后台把日志合并进数据文件。要关注的是压缩节奏,参数compaction.delta_commits表示累积多少个提交后触发一次压缩,默认是 5。如果你的写入频率很高,比如每分钟一次提交,那 5 个提交就是 5 分钟压缩一次,频率偏高会占用资源;设成 10 到 20 更合适。
Clustering 解决的是小文件问题。它会把多个小文件按某种排序规则重组成大文件,既减少文件数,又能提升带条件查询的性能,因为同类的数据被聚到一起了,扫描时可以跳过很多文件。这个操作是异步的,可以在写入作业里配置周期性触发。
-- 在 Hudi 表的 WITH 参数中追加 'clustering.async.enabled' = 'true', 'clustering.schedule.enabled' = 'true', 'clustering.delta_commits' = '20', 'clustering.tasks' = '2'还有两个参数控制文件大小,也很实用:hoodie.parquet.small.file.limit定义什么算小文件,默认 100 MB;hoodie.parquet.max.file.size定义单文件上限,默认 120 MB。调这两个值就能控制治理的激进程度。
实操心得:文件治理不要等到出问题才做。我的做法是作业上线时就把 Clustering 配好,然后每周看一眼分区的平均文件大小,如果持续低于 50 MB,说明配置需要调,别等到查询变慢再回头收拾。
5. 下游对接:TiDB 与 Doris 在链路里的位置
Hudi 解决了存储问题,但它的查询能力偏弱,面对复杂的多维分析还是吃力。所以真实链路里,Hudi 后面通常还要接一层,要么接支持明细查询的数据库,要么接分析型引擎。TiDB 和 Doris 是两条比较常见的路线。
5.1 TiDB + Flink SQL 做明细服务层
TiDB 的定位是支持事务的分布式数据库,水平扩展、兼容常见的 MySQL 协议。在链路里,它适合承担面向在线业务的明细服务层。
举个例子:用户在小程序里点开订单详情,需要毫秒级返回,同时后台运营要按条件筛选订单,这些场景 Hudi 都扛不住。做法是用 Flink SQL 把 Hudi 里的结果再同步到 TiDB 的明细表,对外提供点查和范围查询能力。
Flink SQL 同步到 TiDB 的写法,其实就是把 JDBC 连接器的地址换成 TiDB 的:
CREATE TABLE order_detail_tidb ( order_id STRING, user_id BIGINT, amount DECIMAL(18,2), order_status STRING, update_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:mysql://tidb01:4000/order_center?useSSL=false&rewriteBatchedStatements=true', 'table-name' = 'order_detail', 'username' = 'sync_user', 'password' = '${secret_value}', 'sink.buffer-flush.max-rows' = '1000', 'sink.buffer-flush.interval' = '2s', 'sink.max-retries' = '5' );这里有个经验点:TiDB 侧的建表语句要提前手工建好,并且主键、索引都规划好,不要让同步作业去自动建表。原因是一旦主键没对齐,写入会退化成插入,重复数据悄无声息地堆起来,等你发现时已经很难清理了。
5.2 Doris 做 OLAP 加速与类型对齐
Doris 走的是另一条路,它是列式存储加 MPP 架构,擅长多维聚合查询。链路里的位置通常在 DWS 或者 ADS 层,承担报表和分析。数据流向一般是 Hudi 到 Doris,或者 TiDB 到 Doris。
从 Flink 写 Doris 的参数长这样:
CREATE TABLE ads_daily_gmv ( dt STRING, channel STRING, gmv DECIMAL(20,2), order_count BIGINT, update_time TIMESTAMP(3) ) WITH ( 'connector' = 'doris', 'fenodes' = 'doris-fe01:8030', 'table.identifier' = 'ads.ads_daily_gmv', 'username' = 'writer', 'password' = '${secret_value}', 'sink.label-prefix' = 'flink_ads_daily_gmv', 'sink.properties.format' = 'json', 'sink.properties.read_json_by_line' = 'true', 'sink.enable-delete' = 'true', 'sink.buffer-flush.interval' = '10s' );Doris 这块最容易被类型问题绊住。前面提到的日期类型冲突就经常出现在这里。原因在于 Doris 对日期类型有自己的一套定义,分 DATE 和 DATEV2,前者精度到天,后者支持更广的范围。如果 Flink 侧声明的类型跟 Doris 侧实际列的类型对不上,写的时候就会在传输层报类型不匹配。
我的应对方式是三步走。先在 Doris 侧查DESC table_name,把所有列的真实类型抄下来;再逐字对照 Flink SQL 建表语句;如果确实需要不同精度,就在 SELECT 里用CAST显式转换,不要指望连接器自动帮你猜。这套流程走下来,类型问题基本一次就解决。
Flink SQL里我习惯把''这种空字符串在主键字段上都转成有意义的默认值,因为 Doris 的分区列和主键列对空值处理跟关系型数据库不一样,空值可能被当作特殊值处理,导致分区分不出来。
5.3 一套可复用的对账思路
链路一旦复杂起来,数据不对是迟早的事。我现在的习惯是,每条新链路都配一套对账,不等出问题才补。
对账分三层。第一层是条数对账,比对源端消息数和目标端记录数,看有没有丢。由于 Hudi 是主键去重的,两边的条数本来就可能因为去重而不等,所以要拿主键去重的结果去比,不能直接比总量。第二层是金额或指标对账,把关键度量字段求和做比对,容忍度设一个合理的小数值,处理掉浮点误差。第三层是抽样对账,随机抽一批主键,逐字段比对两边的记录内容,能发现类型转换错误、默认值填充错误这类细节问题。
对账任务我用 Flink SQL 写,把两边的结果都查出来再 FULL JOIN,不一致的直接输出到一张告警表,接告警系统。这个做法比写脚本跑定时任务更实时,出问题能马上知道。
6. Flink 数据血缘:作业跑起来之后才想到的事
作业上线多了,一个绕不开的问题是:某个指标的口径变了,或者某张表的数据有问题,我得知道它影响了多少下游。这就是血缘要解决的事。
6.1 血缘到底要解决什么问题
血缘不是一个“锦上添花”的功能,它解决的是三个很具体的痛点。
第一是影响分析。上游一张表要改结构,我需要立刻知道有哪些作业、哪些报表会受影响。没有血缘的时候,只能靠人工翻代码和文档,漏一个就是线上故障。
第二是问题定位。某个看板数不对,从看板倒推数据源,一路排查是哪一层的过滤条件出了问题。有血缘的话,画一条链路图,顺着往上找就行。
第三是成本归因。一个作业吃掉大量资源,它到底算的是谁的数据、服务的是哪个业务,需要对得上账,否则资源预算没法分配。
6.2 几种落地方式与取舍
血缘采集有几种常见方式,各有适用场景。
一是解析 SQL 语句。把 Flink 作业的 SQL 文本解析成抽象语法树,从中提取表级和字段级的依赖关系。好处是不需要改作业代码,解析器独立运行就行。难点是 SQL 里如果有动态拼接、有自定义函数,解析会不准。
二是借助 SQL 网关。所有作业提交都走一个统一入口,网关在提交时顺便把 SQL 解析出来,把血缘存下来。这种方式比较干净,但前提是团队得接受“所有作业都走网关”这个约束。
三是运行时埋点。在作业里接一些指标上报能力,运行时采集算子之间的数据流向。这种方式最准,但侵入性最强,改造现有作业的成本高。
我自己的选择是先做 SQL 解析,覆盖百分之九十的场景,剩下的复杂情况手工标注。原因很实际:改造成本最低,见效最快,别一上来就追求完美。
6.3 我在实践里的做法
具体落地时,我把血缘信息存成三张表:表级依赖、字段级依赖、作业元信息。表级依赖记录“作业 A 读表 X 写表 Y”,字段级依赖记录“Y 的某列来自 X 的哪几列经过了什么转换”,作业元信息记录负责人、调度周期、资源占用。
采集的触发点放在作业提交环节,解析出结果后写入元数据库。查询时给两个入口:一个是正向的,输入一张表,列出所有下游;一个是反向的,输入一张表,列出所有上游。再配一个图形化展示,把链路画出来,排障时非常直观。
注意:血缘信息会随时间变化,作业改了 SQL,之前解析的依赖就过期了。所以每次作业重新提交都要覆盖更新,而不是追加。我见过有人只追加不更新,最后血缘图里全是历史遗留的假依赖,反而误导排查。
7. 常见问题与排查速查表
这一节我把实际运维中反复遇到的问题整理成速查表,按类型分组,方便直接对照。
7.1 启动与依赖类
| 现象 | 常见根因 | 处理方式 |
|---|---|---|
| 提交作业报 ClassNotFound | 依赖未放入 lib 目录 | 补齐 Hudi bundle、JDBC 驱动、存储客户端 |
| No FileSystem for scheme | 缺少存储客户端或路径前缀写错 | 检查路径前缀与 jar 是否匹配 |
| 版本冲突 NoSuchMethodError | Hudi 包与 Flink 版本不匹配 | 严格按版本矩阵重新下载 |
| 作业启动即 OOM | JobManager 内存过小或依赖过多 | 提高 process.size,精简 lib |
| 无法连接资源管理器 | 配置文件地址或认证有误 | 核对地址、端口、认证配置 |
启动类问题占我遇到问题总量的一半以上,而其中又有八成是依赖问题。所以出现任何“莫名其妙”的报错,我的第一反应永远是去看lib目录里到底有哪些 jar,有没有重复版本,有没有缺。
7.2 写入与存储类
| 现象 | 常见根因 | 处理方式 |
|---|---|---|
| 写入延迟持续升高 | 未开启异步压缩,日志堆积 | 打开 compaction 与 clustering |
| 检查点频繁超时 | 状态过大或存储写入慢 | 开启增量检查点,检查存储带宽 |
| 小文件数量暴涨 | 桶数与数据量不匹配 | 调整桶数,开启 clustering |
| 提交冲突频繁 | 多作业写同一张表 | 合并写入作业,或错开提交时间 |
| 写入吞吐上不去 | write.tasks 太小 | 提高写入并行度,同时核对桶数 |
这里有一个我踩过的坑值得单独说:多个作业同时写同一张 Hudi 表,会频繁发生提交冲突,表现为作业不断重试甚至失败。Hudi 的提交机制需要保证同一时刻只有一个写入者在提交,多写者会互相抢占。解决办法是把写入合并到一个作业里,如果业务上确实需要分开,就通过时间窗口错开,或者使用支持并发写的表类型。
7.3 数据正确性类
| 现象 | 常见根因 | 处理方式 |
|---|---|---|
| 新旧数据颠倒 | precombine 字段选错 | 换成更新时间字段 |
| 目标端条数偏多 | 主键未声明或未对齐 | 补 NOT ENFORCED 主键并核对下游主键 |
| 部分字段为空 | 类型映射不匹配被丢弃 | 用 CAST 显式转换 |
| 日期字段写入失败 | 日期类型语义不一致 | 统一两边类型声明 |
| 迟到数据丢失 | 水位线设置过小 | 调大乱序容忍时间 |
数据正确性问题最麻烦的地方在于它不报错,只是数据悄悄不对。所以我强烈建议把对账做成常态化的,而不是出事了才去查。对账任务的成本很低,但能省下的排查时间是以天计的。
8. 几个我认为最值钱的经验点
学这套组合的过程中,有些经验是文档里不会写、但实际用起来特别关键的。这一节我挑几个最值钱的分享一下。
8.1 参数不要照抄,要理解取值来源
网上能搜到大量 Hudi 和 Flink 的配置模板,直接抄确实能跑起来,但一旦数据量变了、业务节奏变了,抄来的参数立刻就不适用了。我的习惯是每引入一个参数,都问自己一句“这个值的依据是什么”。比如桶数,依据是数据总量除以目标单桶大小;比如检查点间隔,依据是能容忍的数据重复处理时长;比如水位线,依据是数据源的乱序程度。
这三个例子背后是同一套方法论:参数取值应该来自你的业务约束和资源约束,而不是别人的博客。博客能告诉你参数存在、参数怎么配,但取多少必须结合自己的情况算。
8.2 监控看什么指标
作业上线只是开始,能稳定运行才叫完成。我平时盯的指标不多,但每一个都很关键。
延迟方面,看消息积压量和端到端延迟。积压量持续上涨说明处理能力跟不上,要么扩并行度,要么优化算子。延迟方面我会设一个告警阈值,超过就通知。
稳定性方面,看检查点成功率、检查点耗时、失败重启次数。检查点频繁失败说明状态太大或者存储慢,重启次数多说明有隐藏的异常在反复触发。
存储方面,看 Hudi 表的分区文件数、平均文件大小、压缩任务待处理数量。文件数持续增长要警惕,压缩 backlog 一直不下降说明压缩资源不够。
这三类指标加起来不到十个,但覆盖了绝大多数故障的早期信号。指标不是越多越好,能看懂、能行动才有价值。
8.3 学习路径上的一个建议
最后说一个学习方法上的体会。这套组合涉及的知识面很宽:分布式计算、分布式存储、SQL 引擎、消息队列、集群运维。如果按教科书顺序一个个啃,很容易在某个环节卡住就失去动力。
我的做法是从一个最小可跑通的链路开始,先让它跑起来,再逐个环节深挖。先用本地文件系统跑通 Flink SQL 写 Hudi,再换成真实存储,再上集群,再接下游。每前进一步只解决一个问题,这样每次都能看到成果,学习节奏不会断。
还有个细节:刚开始学的时候,别急着上生产规模的数据。用几千条造出来的测试数据跑通流程,比拿几亿条真实数据调参高效得多。等链路逻辑都通了,再逐步加压测边界,这样既安全又省时间。
这套链路我前后折腾了不少时间,中间踩的坑远不止文里写的这些,但核心逻辑其实就那么几条:想清楚数据模型,对齐两端类型,把参数依据算出来,然后坚持做对账。剩下的都是围绕这几条打补丁。