Spark离线数仓与Flink实时数仓双引擎落地:分层架构与实战调优

发布时间:2026/10/7 13:45:18
Spark离线数仓与Flink实时数仓双引擎落地:分层架构与实战调优 简介这是一套面向 Spark 离线数仓与 Flink 实时数仓实战场景的完整项目资料包适合大数据开发、数仓工程师及准备数仓方向面试的读者。压缩包共 607 个文件大小约 54.21MB涵盖大量 Java 源码与 class 编译产物、XML 配置、Shell 部署脚本、SQL 建表脚本、JSON 配置和 Markdown 笔记其中 Java 代码对应数仓各层业务逻辑Shell/SQL 负责环境初始化与任务调度Markdown 则记录分层设计与部署步骤目录结构清晰便于按模块查阅。资料内完整呈现 ODS→DIM/DWD→DWS 的实时数仓分层方案并通过对 Kafka、HBase、ClickHouse、Redis、ES、MySQL 等组件的对比解释了各层选型原因例如使用 HBase 保存维度数据以支持主键查询使用 ClickHouse 承担 DWS 层分析查询离线部分则沉淀了 Spark 数仓的项目实现与部署要点方便同时掌握两类数仓的架构差异与落地方法。目前已有 405 人学习浏览适合作为离线实时一体化数仓建设或面试复盘时的参考资料。1. 这份双引擎数仓源码包里真正值钱的不是代码而是分层拿到「Spark离线数仓Flink实时数仓项目源码部署资料.rar」先别急着解压看代码。类似项目最容易给人的错觉是以为难点在 Spark 或 Flink 的 API 上实际上这套东西真正值钱的是它把“离线一批、实时一条”这两条链路用同一套 ODS/DWD/DWS/ADS 分层规范串了起来。离线由 Spark SQL 跑 T1 批任务实时由 Flink 消费 Binlog 做秒级加工两条线最终在报表层对齐。文章会把这套落地路径拆开讲选型逻辑、分层映射、离线实时怎么写、部署参数怎么调、坑在哪都覆盖。适合三类人准备转数仓开发、公司要搭双跑数仓、以及在 Spark 和 Flink 之间摇摆的从业者。2. 选型先想清楚离线为什么押注 Spark SQL实时为什么押注 Flink2.1 离线数仓为什么选 Spark SQL不是“快”是这三件事很多团队换 Spark 的理由就是“跑得快”但真正让 Spark SQL 替代 Hive on MR 的是三个更实在的点。第一Spark 3.0 之后的 Adaptive Query Execution 会在运行中动态调整 shuffle 分区数对常见的 join 倾斜有一定自愈能力这在写几十张加工表时能省下大量调优时间。第二DataSource v2 和谓词下推让 Spark 读 Hive 表时能把过滤条件下推到文件层而不是全表扫进来再过滤。第三Spark 是统一批处理引擎同一个 SQL 脚本既能跑日批也能应急跑小时级调度不用维护两套代码。但“快”是有前提的。如果你的 executo r内存给得少、shuffle 分区数设置不合理Spark SQL 跑起来比 Hive 还慢的情况我也见过。在这类源码项目里离线链路通常全部用 Spark SQL 加工基本不写 RDD 代码维护成本比老式的 Java MR 项目低一个量级。如果你公司现有离线数仓是 Hive要不要迁 Spark我的判断标准是两条一是有没有大量复杂 join 和窗口函数二是 Yarn 集群资源是否稳定。都满足的话小时级任务提速非常明显否则继续用 Hive 也不算错没必要为了技术热度买单。2.2 实时数仓为什么选 Flink不是“流”是精确一次实时数仓开发工作内容里最难跟新人讲清楚的就是“为什么实时链路不用 Spark Streaming”。Spark Streaming 本质是微批延迟能做到秒级已经很不错但它对事件时间、状态管理和端到端一致性的支持都相对弱。Flink 的 Watermark、State 和两阶段提交让作业既能处理乱序迟到数据又能在重启后做到不丢不重。对实时数仓来说这一条比“吞吐量高”重要得多。这套源码里实时链路的常规组成是MySQL Binlog → Flink CDC → Kafka → Flink SQL 清洗/维表 Join → ClickHouse。Flink 在其中承担的是“同步 清洗 轻度聚合”的活与 Spark 离线链路完全解耦。你会在实时目录里看到大量 CREATE TABLE WITH(connector...) 的语句这就是 Flink SQL 建表的方式。它的优点是上手快缺点是一旦 Connector 参数配错报错信息非常隐晦后面第 5 章会集中讲几个高频异常。2.3 两套分层如何对齐ODS→DWD→DWS→ADS 在 Spark 和 Flink 里的映射两套引擎的分层逻辑必须对齐否则实时报表和离线报表永远对不上账。离线链路里ODS 是 Hive 外部表直接指向 HDFS 上的原始日志目录DWD 做清洗、脱敏和维度退化DWS 做轻聚合ADS 是面向报表导出的宽表。实时链路则对应为ODS 是 Kafka TopicDWD 是清洗/Join 后写回 Kafka 的明细流DWS 是 ClickHouse 里的明细聚合表ADS 是 ClickHouse 对外查询的结果表。-- 离线数仓最小分库脚本按 ODS/DWD/DWS/ADS 四层建库 CREATE DATABASE IF NOT EXISTS ods_xxx COMMENT 原始层 LOCATION /warehouse/ods_xxx; CREATE DATABASE IF NOT EXISTS dwd_xxx COMMENT 清洗层 LOCATION /warehouse/dwd_xxx; CREATE DATABASE IF NOT EXISTS dws_xxx COMMENT 汇总层 LOCATION /warehouse/dws_xxx; CREATE DATABASE IF NOT EXISTS ads_xxx COMMENT 应用层 LOCATION /warehouse/ads_xxx; -- ODS 层用外部表指向原始数据目录数据不移动 CREATE EXTERNAL TABLE IF NOT EXISTS ods_xxx.order_log ( order_id BIGINT, user_id BIGINT, sku_id BIGINT, amount DECIMAL(10,2), status STRING, ts STRING ) PARTITIONED BY (dt STRING) STORED AS PARQUET LOCATION /data/ods/order_log;这里有两个关键设计。一是 ODS 用外部表删除表不影响 HDFS 原始数据生产环境误操作有后悔药二是分区字段 dt 用字符串格式所有离线任务按 dt 调度和回溯实时链路则用事件时间字段对应。两套分层对齐的核心不只是表名一致更关键的是主键、统计口径和时间字段必须同源。否则离线按业务时间统计、实时按处理时间统计结果永远差一截。3. 离线链路落地用 Spark SQL 把 ODS 清洗到 ADS 的完整 SQL 模板3.1 从 ODS 到 DWD清洗、脱敏和维度退化在一段 SQL 里完成离线链路里最耗时的往往不是写 SQL而是写“能直接扔给调度跑的 SQL”。ODS 到 DWD 这一段要做的事很明确过滤掉无效数据、统一字段类型、必要时对敏感字段脱敏、把需要关联的维度退化进来。常见做法是把所有加工写成 INSERT OVERWRITE按天分区重跑保证任务可回溯。SET spark.sql.shuffle.partitions200; SET spark.sql.adaptive.enabledtrue; INSERT OVERWRITE TABLE dwd_xxx.order_dwd PARTITION (dt2024-06-01) SELECT order_id, user_id, COALESCE(sku_id, 0) AS sku_id, CAST(amount AS DECIMAL(10,2)) AS amount, CASE WHEN status IN (paid,shipped,completed) THEN status ELSE unknown END AS status, from_unixtime(CAST(ts AS BIGINT), yyyy-MM-dd HH:mm:ss) AS event_time FROM ods_xxx.order_log WHERE dt 2024-06-01 AND order_id IS NOT NULL;前面两条 SET 不是摆设。spark.sql.shuffle.partitions 控制的是 shuffle 产生的分区数200 是中小数据量下的常用起点数据量翻一个量级后要同步加大。spark.sql.adaptive.enabled 开启后Spark 可以在运行时把过小的分区合并掉减少小文件问题。SQL 本身要遵循一个习惯过滤条件下推能 WHERE 就别 SELECT 后再筛维表字段能退化就提前退化避免下游每层都重复 join 同一张维表。脱敏一般对手机号、身份证这类字段做 md5 或者保留前后几位具体规则按业务定。3.2 从 DWD 到 DWS轻聚合放这层重聚合放 ADS别揉在一起DWD 到 DWS 的这一步最容易犯的错是“把所有聚合都堆到一张表”。轻聚合应该只做细粒度汇总比如按 sku、按小时、按渠道这类常用维度组合而 ADS 层再基于 DWS 做重聚合和宽表拼接。这样做的目的很实际DWS 能被多个 ADS 复用不用每张报表都从明细层重新扫一遍。INSERT OVERWRITE TABLE dws_xxx.order_dws PARTITION (dt2024-06-01) SELECT sku_id, COUNT(order_id) AS order_cnt, SUM(amount) AS gmv, COUNT(DISTINCT user_id) AS uv FROM dwd_xxx.order_dwd WHERE dt 2024-06-01 GROUP BY sku_id;这一段逻辑不复杂但参数上有个经验值得说。COUNT(DISTINCT user_id) 在用户量大的场景是典型的性能杀手如果 UV 精度要求没那么高可以用 approx_count_distinct 替代如果必须精确则要考虑把 UV 明细单独成表而不是每次都从大明细表上 COUNT DISTINCT。另外DWS 层分区策略建议与 DWD 保持一致都是按天分区这样调度依赖和回溯都简单。ADS 层再做一次 GROUP BY 或 JOIN 维表生成宽表给报表系统查询。3.3 调度与依赖比写 SQL 更重要的壳离线数仓跑批最怕的不是 SQL 慢而是下游在空表或脏分区上跑出“全 0 结果”还不报错。我见过太多 Spark 数据分析案例翻车最后定位到是上游 ODS 分区没产出下游照跑不误。所以调度脚本里必须做分区就绪检查分区不存在就直接失败让调度系统重试或告警。#!/bin/bash # 检查上游 Hive 分区是否产出产出才跑本层避免空跑 partitiondt$(date -d yesterday %F) hive -e MSCK REPAIR TABLE ods_xxx.order_log; if hive -e SHOW PARTITIONS ods_xxx.order_log | grep -q $partition; then spark-submit \ --class com.xxx.OfflineJob \ offline-job.jar --date $partition else echo 上游分区未就绪: $partition exit 1 fi这个脚本是离线调度的最小骨架。MSCK REPAIR 是为了让 Hive metastore 识别 HDFS 上新写入的分区尤其当数据是由其他流程直接丢到目录下时分区就绪检查要放在 spark-submit 之前。还有一个细节exit 1 必须在分区缺失时立刻返回否则调度系统会认为任务成功下游一路跑下去最后报表异常时排查成本极高。同样的检查逻辑在 DWS 跑 ADS 之前也要来一道保证整条链路的依赖是显式的。4. 实时链路落地Flink CDC 进 Kafka、维表 Join 与 ClickHouse 攒批4.1 从 MySQL 到 KafkaFlink CDC 建实时 ODS实时链路的 ODS 层最常见做法是用 Flink CDC 直接把 MySQL 业务表同步到 Kafka用 Debezium 格式保留 Binlog 里的 before、after 和 op 字段。这样下游不仅能拿到最新数据还能知道这条数据是插入、更新还是删除。第一次上线时scan.startup.mode 要用 initial它的语义是“先做一次全量快照再无缝切到 Binlog 增量”不会丢数据。-- 实时 ODSMySQL 表通过 CDC 进入 Kafka CREATE TABLE ods_order_mysql ( order_id BIGINT PRIMARY KEY NOT ENFORCED, user_id BIGINT, sku_id BIGINT, amount DECIMAL(10,2), status STRING, ts TIMESTAMP(3) ) WITH ( connector mysql-cdc, hostname mysql-host, port 3306, username cdc_user, password xxx, database-name shop, table-name order, scan.startup.mode initial ); CREATE TABLE kafka_order_sink ( order_id BIGINT PRIMARY KEY NOT ENFORCED, user_id BIGINT, sku_id BIGINT, amount DECIMAL(10,2), status STRING, ts TIMESTAMP(3) ) WITH ( connector kafka, topic ods_order, properties.bootstrap.servers kafka1:9092,kafka2:9092, format debezium-json, sink.semantic exactly-once ); INSERT INTO kafka_order_sink SELECT * FROM ods_order_mysql;这里要提醒两个很现实的坑。第一多个 Flink CDC 作业共用一个 MySQL 实例时必须给每个作业单独设置 server-id否则会跟 Binlog 拉取冲突报错多为“连接被重置”。第二用 Debezium 格式时下游 Flink SQL 必须显式声明主键否则更新和删除事件无法正确路由。如果你只是做 MySQL 到 ClickHouse 的简单同步可以不用 Kafka 中转但只要下游有多个消费者中间放一层 Kafka 几乎是必须的。4.2 实时 DWD维表 Join、脏数据过滤和迟到修正实时 DWD 层的核心操作是维表 Join。Flink SQL 里用 LOOKUP Join 实现语法上要写 FOR SYSTEM_TIME AS OF表示每条流数据到达时去查一次维度表当前版本。维表数据量不大时建议把缓存开大能显著降低对 MySQL 的查询压力数据量大或者更新频繁时则要结合 CDC 维护维度表而不是每次实时查库。CREATE TABLE dim_sku ( sku_id BIGINT PRIMARY KEY NOT ENFORCED, sku_name STRING, category_id BIGINT ) WITH ( connector jdbc, url jdbc:mysql://mysql-host:3306/shop, table-name dim_sku, lookup.cache.max-rows 10000, lookup.cache.ttl 1h, lookup.max-retries 3 ); INSERT INTO kafka_dwd_order SELECT o.order_id, o.user_id, d.sku_name, d.category_id, o.amount, o.ts FROM kafka_order_sink o LEFT JOIN dim_sku FOR SYSTEM_TIME AS OF o.proc_time AS d ON o.sku_id d.sku_id WHERE o.order_id IS NOT NULL;这段 SQL 里WHERE 条件里的 order_id IS NOT NULL 就是实时链路的“分层过滤”。很多新手会省略这一步结果脏数据一路冲到 ClickHouse等报表对账对不上时才发现源头没卡。lookup.cache.ttl 设 1h 表示维度数据在缓存里最多存活一小时适合变化不频繁的维度如果维度每天变好几次要把 ttl 调小到分钟级否则 Join 到的是过期维度。至于迟到数据修正要靠下游窗口计算里的事件时间语义兜底加工时统一用 ts 而不是 proc_time。4.3 实时 DWS/ADS分组聚合之后攒批写 ClickHousesink 参数怎么给实时 DWS 一般落在 ClickHouse。ClickHouse 的写入特性决定了它不适合逐条 Insert高频小写入会产生大量 parts最终触发 “Too many parts” 异常。所以 Flink 写 ClickHouse 一定要开攒批常见的攒批参数是 buffer-flush.max-rows 和 buffer-flush.interval两者满足其一就刷一批。# Flink ClickHouse Sink 攒批参数 sink.buffer-flush.max-rows1000 sink.buffer-flush.interval5s sink.max-retries3CREATE TABLE ch_dws_order ( sku_id BIGINT, order_cnt BIGINT, gmv DECIMAL(10,2), window_start TIMESTAMP(3), PRIMARY KEY (sku_id, window_start) NOT ENFORCED ) WITH ( connector clickhouse, url clickhouse://ch-01:8123, table-name dws_order, sink.buffer-flush.max-rows 1000, sink.buffer-flush.interval 5s ); INSERT INTO ch_dws_order SELECT sku_id, COUNT(order_id) AS order_cnt, SUM(amount) AS gmv, TUMBLE_START(ts, INTERVAL 5 MINUTE) AS window_start FROM kafka_dwd_order GROUP BY sku_id, TUMBLE(ts, INTERVAL 5 MINUTE);攒批的两个参数要按写入吞吐调。1000 条或 5 秒先到先刷是常见起点如果 ClickHouse 压力大把 max-rows 调到 5000、interval 调到 10s 能明显降低 parts 数但报表延迟会略增。窗口聚合这里必须用 TUMBLE(ts, ...) 而不是处理时间否则上游数据一旦延迟重发报表数值就会出现先多后少再修正的“来回跳”这种问题在实时数仓里极其难查尽量从源头避免。另外写入 ClickHouse 建议写本地表再依赖分布式表查询直接大批量写分布式表容易造成节点间数据二次转发。5. 部署与排查把这五个高频坑先填平再动你的集群5.1 spark-submit 参数集群资源与并行度的匹配Spark 集群搭建完成后建议先跑一个简单的 group by 任务确认 shuffle 正常再上业务 SQL。提交离线任务时最常被问的参数就是 executor 个数、内存和 shuffle 分区数。给一个我常用的模板spark-submit \ --master yarn \ --deploy-mode cluster \ --num-executors 20 \ --executor-cores 4 \ --executor-memory 8g \ --conf spark.sql.shuffle.partitions200 \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.sql.adaptive.coalescePartitions.enabledtrue \ --class com.xxx.OfflineJob \ offline-job.jar --date $(date -d yesterday %F)executor-memory 8g、executor-cores 4 是一个平衡点内存再往上容易触发长 GC导致任务看着在跑、进度就是不动内存再往下shuffle 数据会大量落盘磁盘 IO 成为瓶颈。spark.sql.shuffle.partitions 的值需要看 shuffle 数据量200 是中小集群的常用起点数据量大时翻倍。一个很容易忽略的参数是 spark.sql.adaptive.coalescePartitions.enabled开启后小分区会自动合并能减少小文件数量这对下游 Hive 查询收益明显。5.2 Flink 提交与 Checkpoint重启不丢数据的三个必调参数实时作业的提交参数里第一优先级是 Checkpoint。如果不配 Checkpoint作业重启后会从 Kafka 最新位点继续消费中间一段数据永远丢了这种翻车几乎每个团队都遇到过。另一个容易被忽略的是 State Backend状态大时用 RocksDB状态小时用 Heap 即可。flink run -d \ -m yarn-cluster \ -p 4 \ -c com.xxx.RealtimeJob \ -D execution.checkpointing.interval60s \ -D state.backendrocksdb \ -D state.checkpoints.dirhdfs:///flink-checkpoints \ -D execution.checkpointing.min-pause30s \ -D restart-strategyfixed-delay \ -D restart-strategy.fixed-delay.delay10s \ realtime-job.jarexecution.checkpointing.interval 设为 60s 是一个稳妥起点太频繁会导致磁盘和网络压力大太稀疏则故障恢复时丢失的数据多。min-pause30s 表示两次 Checkpoint 之间至少间隔 30 秒避免一次还没做完下一次又启动。RocksDB 状态后端在算子状态超过几百 MB 时优势明显但它有本地磁盘依赖容器化部署时要给 TaskManager 挂可靠的本地盘。重启策略用 fixed-delay适合大部分实时链路如果上游 binlog 延迟导致作业反复重启要配合告警而不是无限重试。5.3 高频坑记录现象、原因与解决先说 Flink JDBC 连接器异常。现象是作业运行一段时间后报“Connection is not available, request timed out”通常出现在维表 Join 或 Sink 到 MySQL 的链路上。原因是并行度太高默认 JDBC 连接池上限不够用大量请求在排队等连接。解决方法是给维表开启缓存降低查询频率或者调大连接池更省事的做法是把维表放进 Redis用异步 IO 查 Redis。第二个坑是使用 Flink 同步 MySQL 到 ClickHouse 时报 “Too many parts”。现象是作业稳定运行很久后某天突然写入失败ClickHouse 日志里全是 parts 超限。原因通常是攒批参数设得过大或分区键设计不合理单分区内 parts 数量超过 merge 速度。解决方法是调小 buffer-flush.interval把写入分散到更细的时间分区并且优先写本地表而不是分布式表。第三个坑是 Spark SQL 数据倾斜。现象是某个 executor 长时间卡住甚至直接 OOM其他 executor 早就跑完了。原因常见于 join 或 group by 的 key 集中在少数热值上比如某个爆款 sku 占据了大部分数据。解决方法是先开 AQE再考虑对热点 key 做加盐处理比如把大 key 随机拆成多份两阶段聚合后再合并结果。第四个坑是实时和离线口径对不上。现象是离线报表和实时大屏的 GMV 永远有差值而且差值不固定。原因大多是实时链路用了处理时间离线链路按业务时间统计数据只要稍有延迟两边归属的日期就不一样。解决方法是统一定义事件时间字段实时窗口和离线 SQL 都基于它计算双跑期每天都跑对账差值控制在阈值内才允许切流。第五个坑是加了 Kafka 分区后 Flink 吞吐上不去。现象是分区从 6 加到 12消费速率几乎没变。原因是 Flink 的并行度大于分区数时多余的并行度是空闲的一个并行度最多消费一个分区。解决方法是让并行度和分区数对齐想提升吞吐优先加分区而不是加并行度。6. 验证进阶流批双跑对账一个让你睡好觉的土办法新链路最怕的不是跑不通而是跑通了但数值没人敢信。实时链路刚上线时我的习惯是保留至少一周的“双跑期”离线照常出 T1 报表实时大屏同步跑每天用一张对账 SQL 拉两边差异。这里用不上复杂工具一条 SQL 就够了。-- 离线 ADS 与实时 ClickHouse 按日对账diff 不为 0 的记录打印出来 SELECT COALESCE(a.stat_date, b.stat_date) AS stat_date, COALESCE(a.sku_id, b.sku_id) AS sku_id, a.gmv AS offline_gmv, b.gmv AS online_gmv, a.gmv - b.gmv AS diff FROM ( SELECT dt AS stat_date, sku_id, SUM(gmv) AS gmv FROM ads_xxx.order_ads WHERE dt 2024-06-02 GROUP BY dt, sku_id ) a FULL OUTER JOIN ( SELECT toDate(window_start) AS stat_date, sku_id, sum(gmv) AS gmv FROM clickhouse_db.dws_order WHERE window_start 2024-06-02 00:00:00 AND window_start 2024-06-03 00:00:00 GROUP BY stat_date, sku_id ) b ON a.sku_id b.sku_id AND a.stat_date b.stat_date WHERE a.gmv - b.gmv 0 ORDER BY diff DESC;这条 SQL 的用途是每天上班先跑一遍而不是出问题之后再查。diff 不为 0 时优先排查三件事窗口时间字段是否一致、维度表是否同一版本、Kafka 是否有积压未消费。由于迟到数据的存在实时结果允许与离线有少量偏差但如果某个 sku 的差值一直存在且方向固定基本可以断定是加工逻辑问题不是数据延迟。我给自己定的规矩是离线与实时必须双跑满一周确认 diff 连续三天在阈值内才允许把实时报表挂到正式大屏上。这个土办法救过我很多次看似占用了额外资源实际上比上线后对账排查省钱得多。希望帮到你。本文还有配套的精品资源点击获取