简介:这是一套面向 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 跑 T+1 批任务,实时由 Flink 消费 Binlog 做秒级加工,两条线最终在报表层对齐。文章会把这套落地路径拆开讲,选型逻辑、分层映射、离线实时怎么写、部署参数怎么调、坑在哪都覆盖。适合三类人:准备转数仓开发、公司要搭双跑数仓、以及在 Spark 和 Flink 之间摇摆的从业者。
2. 选型先想清楚:离线为什么押注 Spark SQL,实时为什么押注 Flink
2.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 Topic,DWD 是清洗/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.partitions=200; SET spark.sql.adaptive.enabled=true; INSERT OVERWRITE TABLE dwd_xxx.order_dwd PARTITION (dt='2024-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 (dt='2024-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 分区是否产出,产出才跑本层,避免空跑 partition="dt=$(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 到 Kafka:Flink CDC 建实时 ODS
实时链路的 ODS 层,最常见做法是用 Flink CDC 直接把 MySQL 业务表同步到 Kafka,用 Debezium 格式保留 Binlog 里的 before、after 和 op 字段。这样下游不仅能拿到最新数据,还能知道这条数据是插入、更新还是删除。第一次上线时,scan.startup.mode 要用 initial,它的语义是“先做一次全量快照,再无缝切到 Binlog 增量”,不会丢数据。
-- 实时 ODS:MySQL 表通过 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:分组聚合之后攒批写 ClickHouse,sink 参数怎么给
实时 DWS 一般落在 ClickHouse。ClickHouse 的写入特性决定了它不适合逐条 Insert,高频小写入会产生大量 parts,最终触发 “Too many parts” 异常。所以 Flink 写 ClickHouse 一定要开攒批,常见的攒批参数是 buffer-flush.max-rows 和 buffer-flush.interval,两者满足其一就刷一批。
# Flink ClickHouse Sink 攒批参数 sink.buffer-flush.max-rows=1000 sink.buffer-flush.interval=5s sink.max-retries=3CREATE 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.partitions=200 \ --conf spark.sql.adaptive.enabled=true \ --conf spark.sql.adaptive.coalescePartitions.enabled=true \ --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.interval=60s \ -D state.backend=rocksdb \ -D state.checkpoints.dir=hdfs:///flink-checkpoints \ -D execution.checkpointing.min-pause=30s \ -D restart-strategy=fixed-delay \ -D restart-strategy.fixed-delay.delay=10s \ realtime-job.jarexecution.checkpointing.interval 设为 60s 是一个稳妥起点,太频繁会导致磁盘和网络压力大,太稀疏则故障恢复时丢失的数据多。min-pause=30s 表示两次 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. 验证进阶:流批双跑对账,一个让你睡好觉的土办法
新链路最怕的不是跑不通,而是跑通了但数值没人敢信。实时链路刚上线时,我的习惯是保留至少一周的“双跑期”,离线照常出 T+1 报表,实时大屏同步跑,每天用一张对账 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 连续三天在阈值内,才允许把实时报表挂到正式大屏上。这个土办法救过我很多次,看似占用了额外资源,实际上比上线后对账排查省钱得多。希望帮到你。
本文还有配套的精品资源,点击获取