☰
数据治理与实时数仓:Flink SQL + Hudi 流批一体方案实践
2026/10/3 5:20:34 网站建设 项目流程

简介:这是一份面向化工园区智能化管控平台的数据处理和存储系统建设方案,适合政府、园区及企业信息化规划人员,用于前期可研、初步设计或软硬件选型参考。文档以系统结构、数据计算、数据存储、数据传输为主线:先按三类用户和300个用户规模框定资源需求;再以TPC-C基准测算数据库服务器峰值处理能力,推算虚拟化服务器数量;随后结合系统数据、业务数据和非结构化数据增量,给出三年8.1TB存储容量配置;最后估算网络带宽并给出设备选型与集群冗余建议。全文含目录、计算公式和服务器配置表,便于直接复用。资源包为单个doc文档,约183KB,结构完整,已有222人学习下载,适合在同类项目容量规划中借鉴。

1. 先说结论:这套 1.0 方案解决的是“数据进了数仓却用不起来”的账

很多团队做数据处理和存储系统建设方案时,第一反应是先买 Hadoop、搭集群、建库建表,结果三个月后发现最痛的不是没数据,而是业务要的实时指标离线给不了、离线要的明细实时又塞不下。这套方案的目标很直接:把数据从产生到可查询拆成接入、清洗、存储、消费四段,用一套统一的数据处理框架串起来,让实时和批量共用同一张表,业务侧只需要关心查询。

谁适合照着做?数据量在几十 TB 到数 PB 之间、既有 T+1 报表又有实时大屏、并且暂时不想上云原生数仓的团队。看完之后你能回答三个问题:存储层该用哪些组件、Flink 作业怎么写不翻车、上线前怎么验证这套系统真的能用。

2. 架构拆分:五个层把“数据处理”和“存储系统”钉死在一张图里

一个方案如果从组件讲起,读者会迷失在“该用 ClickHouse 还是 Doris”“该上 Hudi 还是 Iceberg”的争论里。我习惯先把五个槽位定死,再往槽位里填组件。

接入层负责把数据从业务系统里搬出来;缓冲层用消息队列削峰,避免业务高峰期把后端的库压垮;计算层承担统一的数据处理框架职责,实时和批量用同一套 SQL;存储层管明细和聚合结果;管理层则管调度、权限、数据质量巡检。

2.1 先定边界:不是所有数据都该进这套系统

最常见的失败原因不是组件选错,而是把什么数据都往里塞。下面是我做方案时固定过一遍的接入边界表:

数据形态接入方式实时性建议落点
业务库明细Flink CDC 或 Debezium秒级ODS 层,保留原始结构
客户端埋点日志日志采集 -> Kafka分钟级ODS 层,按事件类型分主题
文件型外部数据定时上传到对象存储T+1独立分区,按日期扫描
强事务账务数据不接入不接入留在业务数据库,只读视图同步

边界之外有两类数据我强烈建议不碰。一类是需要强事务和行级锁的账务核心表,一旦消息乱序或者任务重放,账就对不上,这个锅数据平台背不起;另一类是只活半年、数据量几百 GB 的临时接口数据,直接在业务库里查更快,为它建管道纯属给自己挖坑。

2.2 数据分层映射:ODS / DWD / DWS / ADS 在存储上的落位

传统数仓的分层概念在这个方案里不是摆设,它直接决定对象存储上的目录结构和表格式选择。ODS 层是原样落地的原始数据,Kafka 里的所有事件按日期分区写进对象存储,这一层基本不承担查询压力,但它是数据重跑和审计的后悔药。

DWD 层是清洗后的明细,以业务主键为准做 upsert,这一层用 Hudi 的 MERGE_ON_READ 表类型。DWS 层放轻度聚合结果,订单数、流水额这类按 5 分钟或天聚合的宽表,落到 Doris 或 ClickHouse 供 BI 高频查询。ADS 层就简单了,直接是视图或者几十行的小表,给大屏和报表用。

分层内容存储与格式更新方式典型查询
ODS原始事件对象存储 + Hudi COPY_ON_WRITE追加写几乎不查,仅审计和重跑
DWD主键明细对象存储 + Hudi MERGE_ON_READ按主键 upsert取数、二次加工
DWS轻度聚合宽表Doris / ClickHouse批次覆写BI、实时大屏
ADS指标结果视图或小表定时计算报表接口

路径命名也必须在方案里定死,我一般用s3a://lake/ods/order/ds=2024-05-01/这种格式,分区字段统一叫ds。注意 Hudi 表不是纯目录,它除了数据文件还有.hoodie元数据目录,千万不要把 ODS 和 DWD 的表目录放在同一个前缀下面,后面做权限隔离时会很难受。

2.3 Lambda 与 Kappa 的选型:为什么我选了流批一体

这是方案评审会上一定会被问的问题。Lambda 架构是流批两套代码,实时链路用 Flink,离线链路用 Spark,各有各的表;Kappa 架构是流批共用一套代码、一张表。

维度LambdaKappa 流批一体
代码维护两套 SQL、两套调度、两套告警一套 Flink SQL 复用
存储一致性流表和批表可能对不上同一张 Hudi 表,时点一致
数据恢复重跑整个离线链路重放 Kafka + 恢复 checkpoint
适合场景已有重型离线数仓、批处理逻辑复杂从零建设、团队规模小

我的结论很直接:从零建设数据处理和存储系统,直接走 Kappa,用同一套 Flink SQL 处理网约车订单这种高频明细数据。跑批需求来了也不用换引擎,Flink 批模式执行同一段 SQL 就行。Kappa 不是万能的,它对 Kafka 的消息保留时长要求很高,方案里我要求 Kafka topic 保留至少 7 天,万一任务出问题还能原地重放。

2.4 资源隔离与调度:批量任务和实时任务别抢同一个队列

流批一体并不代表所有任务混在一个资源池里。我把 Yarn 队列按容量拆成两个,实时作业单独占三成,批量作业占七成,避免凌晨批量任务把队列塞满,导致实时作业的 checkpoint 持续延迟,延迟几下之后 Flink 任务就会自动 failover,这是很隐蔽的故障。

# capacity-scheduler.xml 中拆分实时与批量队列 yarn.scheduler.capacity.root.realtime.capacity=30 yarn.scheduler.capacity.root.batch.capacity=70 yarn.scheduler.capacity.root.realtime.maximum-capacity=40 yarn.scheduler.capacity.root.batch.maximum-capacity=80

参数说明:capacity是队列保证的绝对容量,maximum-capacity是队列能借到的上限。实时队列上限设 40 而不是 100,是为了防止批量任务空闲时实时任务把整个集群吃满,等批量任务回来时反而抢不到资源。实际运维中我给实时作业的 Flink 并行度按 Kafka 分区数来定,这个细节第 4 章会专门讲。

3. 存储层先落地:为什么选“对象存储 + Hudi”,参数怎么设

数据处理的底座是存储,存储选型的错误会在三个月后集中爆发。这里说的存储不是单指 HDFS,而是整个存储层方案。我的选择是对象存储做主存储,Hudi 做表格式,HDFS 只留少量路径跑历史遗留的 Spark 任务。

3.1 存储选型:对象存储和 HDFS 的真实差距

HDFS 在小文件场景下非常痛苦,NameNode 的内存有限,几百万个几 KB 的小文件直接让元数据服务变慢;对象存储则没有这个压力,文件即对象,目录只是逻辑概念。HDFS 的另一个问题是扩容要动硬件,对象存储无论是自建 MinIO、Ceph 还是直接用云上的对象存储服务,扩容都是横向加节点或者直接提配额。

对比项HDFS对象存储
小文件元数据压力大,需合并无压力,按对象存储
扩容加节点、做均衡横向扩展,成本低
语义强一致 rename部分对象存储最终一致,要注意
与 Hudi 配合支持好支持好,推荐 S3A 协议

注意对象存储的最终一致性是个坑。自建 Ceph 在并发写同路径文件时可能出现短时间读到旧对象的情况,所以 Hudi 表的写入路径要避免多个作业同时写同一个分区,这点落实到调度上就是同一张表只允许一个 Flink 作业写,其他作业只读。

3.2 Hudi 表的关键参数:主键、预组合字段和小文件治理

Hudi 表不是建完就完事的,参数设置直接决定它是帮你省心还是给你添乱。最有价值的几个参数我来逐个说清楚。

参数推荐值作用与踩坑
hoodie.datasource.write.recordkey.field业务主键不设置或设置错误,upsert 退化成 append
hoodie.datasource.write.precombine.field业务时间戳用摄入时间会导致迟到数据覆盖正确结果
hoodie.parquet.small.file.limit134217728小于该值的文件会被合并,默认也够用
hoodie.clustering.inlinetrue在线聚簇,配合inline.max.commits=4
hoodie.archive.commits.retained20保留最近提交记录数,太小没法时间旅行

recordkey.field是 Hudi 判断主键的唯一依据,订单表就是order_id。precombine.field我强调过很多次,必须用业务时间戳,也就是订单发生时间ts,而不是 Flink 的处理时间。原因很简单:如果一条迟到的订单数据晚到了 10 分钟,它的业务时间更早,按业务时间才能正确决定新旧,按摄入时间则会让这条迟到的旧数据把正确的新数据覆盖掉。

3.3 建表落地:用 Flink SQL 创建 DWD 表,参数全注释

参数说再多不如直接看一张建表语句。下面这张是订单明细的 DWD 表,也是整套方案里最核心的一张表。

CREATE TABLE dwd_order_hudi ( order_id BIGINT, driver_id BIGINT, passenger_id BIGINT, start_lng DOUBLE, start_lat DOUBLE, end_lng DOUBLE, end_lat DOUBLE, amount DECIMAL(10,2), status STRING, ts TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'connector' = 'hudi', 'path' = 's3a://data-lake/dwd/order', 'table.type' = 'MERGE_ON_READ', 'write.operation' = 'upsert', 'hoodie.datasource.write.recordkey.field' = 'order_id', 'hoodie.datasource.write.precombine.field' = 'ts', 'hoodie.parquet.small.file.limit' = '134217728', 'hoodie.clustering.inline' = 'true', 'hoodie.clustering.inline.max.commits' = '4' );

表类型选MERGE_ON_READ是因为订单明细写入频繁,COW 表每次 upsert 都要重写整个文件,写入放大会让高峰期作业背不住;MOR 把更新先记到 log 文件里,查询时再合并,读性能略差但写路径稳。PRIMARY KEY (order_id) NOT ENFORCED是 Flink 连接器的写法,Hudi 不强制在引擎层做唯一性校验,真正的主键约束由 Hudi 的 record key 机制负责。

4. 用 Flink SQL 把流和批收进同一张表:最小可跑链路

架构定完,存储层建完表,接下来就是把数据处理链路跑起来。我拿网约车订单数据这个大家最熟悉的场景举例,从 Kafka 接入原始订单事件,清洗掉缺失值和异常值,写入 Hudi 表。整条链路只用 Flink SQL,不用写一行 Java。

4.1 目录与依赖:一套最小可跑的 Flink SQL 工程

主线作业不需要复杂工程结构,目录干净点反而好维护。我一般建这么几个目录,sql 目录放 DDL 和 INSERT 语句,scripts 目录放提交脚本,checkpoints 目录放本地验证时的状态文件。

mkdir -p /opt/data-platform/{conf,sql,scripts,checkpoints}

以下是提交 SQL 作业的入口脚本,我用它统一管理所有 Flink SQL 任务,避免每个人记住一长串参数。

#!/usr/bin/env bash FLINK_HOME=/opt/flink SQL_FILE=${1:-sql/dwd_order.sql} "$FLINK_HOME"/bin/flink-sql-client.sh \ -D execution.checkpointing.interval=60000 \ -D state.backend=filesystem \ -D state.checkpoints.dir=s3a://data-lake/checkpoints \ -D parallelism.default=8 \ -f "$SQL_FILE"

脚本说明:execution.checkpointing.interval是 checkpoint 间隔,生产环境我至少设 60 秒,太频繁会让对象存储的写入压力变大;state.backend=filesystem把 Flink 状态放到文件系统,配合state.checkpoints.dir指定的路径,任务挂了才能恢复。parallelism.default=8是默认并行度,具体值要根据 Kafka 分区数来定,这个下文细说。

4.2 建 Kafka 源表和 Hudi 目标表:缺失值、异常值在哪一步处理

源表定义直接对应 Kafka 里的原始订单事件。这里有个细节:JSON 里如果混入了脏字段,json.ignore-parse-errors打开后解析失败的行会被丢掉,不会让整个作业卡死。

CREATE TABLE ods_order_mq ( order_id BIGINT, driver_id BIGINT, passenger_id BIGINT, start_lng DOUBLE, start_lat DOUBLE, end_lng DOUBLE, end_lat DOUBLE, amount DECIMAL(10,2), status STRING, ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL '30' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'ods_order_mq', 'properties.bootstrap.servers' = 'kafka1:9092,kafka2:9092', 'properties.group.id' = 'ods_order_consumer', 'format' = 'json', 'json.ignore-parse-errors' = 'true' );

WATERMARK 语句声明了事件时间的乱序容忍度为 30 秒,超过这个时间还没到的数据会被当作迟到数据丢弃。如果你公司的业务场景里订单回调经常延迟几分钟,就要把 30 秒调大,但同时要清楚:容忍度越大,窗口计算的结果出来越晚。

目标表沿用第 3 章的dwd_order_hudi,DROP 掉重建也行,生产上一旦有数据就别轻易 DROP。清洗逻辑全部放在 INSERT 语句里,不在 DDL 里做。

4.3 写一个能反复重放的 INSERT:清洗逻辑与幂等保障

下面的 INSERT 语句把数据从 Kafka 源表写入 Hudi 目标表,顺带把这套方案的清洗规则完整展示出来。

INSERT INTO dwd_order_hudi SELECT order_id, driver_id, COALESCE(passenger_id, 0) AS passenger_id, COALESCE(start_lng, end_lng) AS start_lng, COALESCE(start_lat, end_lat) AS start_lat, end_lng, end_lat, CAST(COALESCE(amount, 0) AS DECIMAL(10,2)) AS amount, status, ts FROM ods_order_mq WHERE status IN ('FINISHED', 'CANCELLED') AND ABS(amount) < 10000 AND start_lng BETWEEN 73 AND 136 AND start_lat BETWEEN 3 AND 54 AND end_lng BETWEEN 73 AND 136 AND end_lat BETWEEN 3 AND 54;

逻辑说明:COALESCE处理缺失值,passenger_id为空时补 0,经纬度为空时用终点经纬度兜底,金额为空时补 0;ABS(amount) < 10000把金额为负和金额异常大的脏数据过滤掉;经纬度范围判断把明显越界的异常值挡在存储层之外。这里说的 dataframe 层面的缺失值和异常值处理习惯,换成 SQL 一样适用,只是处理的对象从单机 DataFrame 变成了流式数据。

这条 INSERT 可以反复重放而不产生脏数据,原因有两个:第一,Flink 的 Kafka 源会定期提交 offset,配合 checkpoint 实现精确一次语义;第二,Hudi 表按order_idupsert,同一订单重复写入时按precombine.field也就是ts判断新旧,旧数据不会覆盖新数据。所以任务失败后直接重启同一个 SQL 文件,消费位置从最近 checkpoint 继续,不需要人工清数。

4.4 必调参数:checkpoint 与并发度的第一版配置

流式入湖的第一个玄学点就是 checkpoint 和并发度,这两个参数翻车率最高。

参数推荐配置说明
execution.checkpointing.interval60000 ms生产最低 60 秒,本地验证可设 10 秒
execution.checkpointing.tolerable-failed-checkpoints3连续失败 3 次才让作业失败,避免抖动
parallelism.default= Kafka 分区数大于分区数时多出的 task 空转,小于时消费能力不够
sink.parallelism与 source 相同Hudi sink 并发太高会同时写大量小文件

关键点是并行度不要拍脑袋设 32。如果 Kafka topic 只有 12 个分区,源并行度设 12 就够,设 32 后多出的 20 个 task 完全空闲,checkpoint 还要等它们确认,反而拖慢整体进度。目标表并行度保持一致,让每个子任务只处理自己负责的主键范围,小文件数量也可控。

5. 五个高频翻车点:现象、原因、解决

方案能不能在线上站稳,看的是这些坑有没有提前填平。我把踩过的坑按“现象 -> 原因 -> 解决”整理成五条,每一条都是真金白银换来的经验。

5.1 对象存储 Token 过期,长任务跑到一半翻车

现象:数据量大的 Flink 作业稳定跑 6 小时后突然报AccessDeniedException,作业自动重启后依旧在同一位置失败,日志里的主键和时间戳都对得上,唯独代码里用的访问凭证失效了。

原因:对象存储的临时凭证有有效期,很多云厂商的临时 Token 默认最长 12 小时,而长任务和重跑任务很容易跨越这个时间点。

解决:接入层统一封装一个凭证刷新组件,作业启动时从鉴权服务换取凭证,并周期性刷新。不想上组件的团队也要做两件事:规划任务时长时预留凭证提前量,生产环境禁止把 AK/SK 明文固化在作业代码里。

5.2 小文件治理忘记开,三个月后表查询和提交一起变慢

现象:Hudi 表刚上线时查询很快,三个月后一张按天分区的订单表产生了上万个几十 KB 的 parquet 文件,Flink 提交一个新的 commit 要扫几十万个文件,查询更是慢到分钟级。

原因:并行度过高、写入频率高,每个 task 各自写各自的小文件,线上没有开 clustering 聚簇,Hudi 自带的small.file.limit合并逻辑也没有触发。

解决:三个参数必须一起开。hoodie.parquet.small.file.limit=134217728控制小于 128MB 的文件参与合并,hoodie.clustering.inline=true配合hoodie.clustering.inline.max.commits=4每 4 个 commit 做一次聚簇。已经产生的表用离线 clustering 任务跑一次,推荐按分区逐个聚簇,避免一次聚太多把对象存储打满。

5.3 并行度大于 Kafka 分区数,Checkpoint 卡着不动

现象:作业状态显示RUNNING,但 checkpoint 一直不成功,Kafka 消费延迟持续增长,业务方反馈大屏指标已经落后半小时。

原因:Flink 的 checkpoint barrier 需要在所有 source task 之间对齐,并行度设为 16 但 Kafka 主题只有 8 个分区,多出的 8 个 task 没有数据可读,barrier 永远无法对齐。

解决:source 并行度严格等于或整数倍于 Kafka 分区数。我一般直接设为相等,因为整数倍也只是让多出的 task 空转。这个检查列入上线 checklist:建主题时定分区数,写作业时抄分区数,两者不一致直接拦截发布。

5.4 预组合字段用错,迟到数据把正确结果覆盖了

现象:某一天订单表里的部分记录,amount字段回到了昨天晚上的值,而实际上今天白天已经更新过正确金额。整个 DWS 聚合结果跟着出错。

原因:预组合字段用了 Flink 处理时间PROCESS_TIME而不是业务时间ts。晚到的历史数据在 Flink 处理时时间更晚,按处理时间判断新旧,这条旧数据反而被判定为“新数据”覆盖了正确值。

解决:hoodie.datasource.write.precombine.field必须指向业务事件时间,也就是订单发生时间。如果业务时间存在但格式不统一,先统一转成TIMESTAMP(3)再写入。这个字段选错,Hudi 表的整表可信度都会被打问号。

5.5 时区没对齐,凌晨高峰期的数据全被算到了零点

现象:大屏上的今日订单数每天 8 点前都异常低,过了 8 点又突然跳涨,检查 Hudi 表发现凌晨 7 点到 8 点的订单时间戳全部变成了当天 0 点。

原因:Flink 解析 JSON 里的时间字符串时按 UTC 处理,没有做时区转换,东八区凌晨的数据落库后全部变成前一天的“深夜”。

解决:Kafka JSON 里的时间字段先用字符串读取,SQL 内显式转换时区。推荐写CONVERT_TZ(ts_str, 'UTC', 'Asia/Shanghai'),不要依赖 Flink 集群默认时区。检查方式很简单:每天看一次max(ts)和当前时间的差值,偏差超过 1 小时就要怀疑时区处理。

6. 上线前怎么验证这套系统真的能用:质量巡检与延迟监控

方案上线最怕的是“能查了”就当成功,实际上数据对不对、延迟高不高完全没数。我的验证习惯分三步:数据质量巡检、链路延迟核对、最小血缘盘点。第一步做一张质量巡检表,用定时 SQL 跑核心表的空值率和主键唯一性。

SELECT COUNT(*) AS total_cnt, COUNT(DISTINCT order_id) AS unique_cnt, SUM(CASE WHEN driver_id IS NULL THEN 1 ELSE 0 END) AS missing_driver_cnt, SUM(CASE WHEN amount < 0 THEN 1 ELSE 0 END) AS negative_amount_cnt FROM dwd_order_hudi WHERE ts >= NOW() - INTERVAL '1' DAY;

巡检的目的不是消灭所有异常,而是给异常设阈值。比如唯一键数量与总行数的比值低于 0.999,说明有重复主键;missing_driver_cnt占比超过 1% 说明上游采集链路有字段丢失。把这些阈值写进告警,比业务方投诉更早发现问题。

链路延迟我用两个指标对照:Kafka consumer lag 和 Hudi 表max(ts)与当前时间的差值。前者看消费端有没有积压,后者看数据从 Kafka 写入 Hudi 的端到端延迟,还可以加一张延迟明细表记录每个分区的最大事件时间,及时发现某个分区卡住。

最后是血缘盘点。我习惯上线前手工把最热的三张表画出血缘图,确认每张表的来源、清洗规则和下游消费方,再让平台自动采集完整的血缘关系。这样每次数据出问题,都能从 ADS 指标一路追到 ODS 原始事件,再决定是重刷还是补数。这套方案走到这一步才算真正立住:处理链路可重放,存储结果可回滚,质量变化可告警。希望帮到你。

本文还有配套的精品资源,点击获取

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

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

立即咨询