ARTICLE DETAIL

资讯详情

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

数据治理与实时数仓:Flink SQL + Hudi 流批一体方案实践

数据治理与实时数仓:Flink SQL + Hudi 流批一体方案实践 简介这是一份面向化工园区智能化管控平台的数据处理和存储系统建设方案适合政府、园区及企业信息化规划人员用于前期可研、初步设计或软硬件选型参考。文档以系统结构、数据计算、数据存储、数据传输为主线先按三类用户和300个用户规模框定资源需求再以TPC-C基准测算数据库服务器峰值处理能力推算虚拟化服务器数量随后结合系统数据、业务数据和非结构化数据增量给出三年8.1TB存储容量配置最后估算网络带宽并给出设备选型与集群冗余建议。全文含目录、计算公式和服务器配置表便于直接复用。资源包为单个doc文档约183KB结构完整已有222人学习下载适合在同类项目容量规划中借鉴。1. 先说结论这套 1.0 方案解决的是“数据进了数仓却用不起来”的账很多团队做数据处理和存储系统建设方案时第一反应是先买 Hadoop、搭集群、建库建表结果三个月后发现最痛的不是没数据而是业务要的实时指标离线给不了、离线要的明细实时又塞不下。这套方案的目标很直接把数据从产生到可查询拆成接入、清洗、存储、消费四段用一套统一的数据处理框架串起来让实时和批量共用同一张表业务侧只需要关心查询。谁适合照着做数据量在几十 TB 到数 PB 之间、既有 T1 报表又有实时大屏、并且暂时不想上云原生数仓的团队。看完之后你能回答三个问题存储层该用哪些组件、Flink 作业怎么写不翻车、上线前怎么验证这套系统真的能用。2. 架构拆分五个层把“数据处理”和“存储系统”钉死在一张图里一个方案如果从组件讲起读者会迷失在“该用 ClickHouse 还是 Doris”“该上 Hudi 还是 Iceberg”的争论里。我习惯先把五个槽位定死再往槽位里填组件。接入层负责把数据从业务系统里搬出来缓冲层用消息队列削峰避免业务高峰期把后端的库压垮计算层承担统一的数据处理框架职责实时和批量用同一套 SQL存储层管明细和聚合结果管理层则管调度、权限、数据质量巡检。2.1 先定边界不是所有数据都该进这套系统最常见的失败原因不是组件选错而是把什么数据都往里塞。下面是我做方案时固定过一遍的接入边界表数据形态接入方式实时性建议落点业务库明细Flink CDC 或 Debezium秒级ODS 层保留原始结构客户端埋点日志日志采集 - Kafka分钟级ODS 层按事件类型分主题文件型外部数据定时上传到对象存储T1独立分区按日期扫描强事务账务数据不接入不接入留在业务数据库只读视图同步边界之外有两类数据我强烈建议不碰。一类是需要强事务和行级锁的账务核心表一旦消息乱序或者任务重放账就对不上这个锅数据平台背不起另一类是只活半年、数据量几百 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/ds2024-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.capacity30 yarn.scheduler.capacity.root.batch.capacity70 yarn.scheduler.capacity.root.realtime.maximum-capacity40 yarn.scheduler.capacity.root.batch.maximum-capacity80参数说明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 退化成 appendhoodie.datasource.write.precombine.field业务时间戳用摄入时间会导致迟到数据覆盖正确结果hoodie.parquet.small.file.limit134217728小于该值的文件会被合并默认也够用hoodie.clustering.inlinetrue在线聚簇配合inline.max.commits4hoodie.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.interval60000 \ -D state.backendfilesystem \ -D state.checkpoints.dirs3a://data-lake/checkpoints \ -D parallelism.default8 \ -f $SQL_FILE脚本说明execution.checkpointing.interval是 checkpoint 间隔生产环境我至少设 60 秒太频繁会让对象存储的写入压力变大state.backendfilesystem把 Flink 状态放到文件系统配合state.checkpoints.dir指定的路径任务挂了才能恢复。parallelism.default8是默认并行度具体值要根据 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_hudiDROP 掉重建也行生产上一旦有数据就别轻易 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经纬度为空时用终点经纬度兜底金额为空时补 0ABS(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.limit134217728控制小于 128MB 的文件参与合并hoodie.clustering.inlinetrue配合hoodie.clustering.inline.max.commits4每 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 原始事件再决定是重刷还是补数。这套方案走到这一步才算真正立住处理链路可重放存储结果可回滚质量变化可告警。希望帮到你。本文还有配套的精品资源点击获取
返回列表