
简介面向大数据离线与实时数仓开发者的完整项目源码与部署资料包涵盖Spark离线数仓与Flink实时数仓两条技术主线。资源围绕实时数仓分层设计展开包括ODS层Kafka接入、DWD层流式处理、DWS层ClickHouse汇总以及DIM层HBase维表存储等核心环节并给出各存储选型的对比分析。压缩包共607个文件以Java源码与class编译产物为主同时包含SQL建表脚本、Shell部署脚本、XML/JSON配置、MD说明文档及项目配置等整体大小54.21MB目录结构清晰便于按模块学习或二次开发。已有405人学习下载。读者可从中获取可直接运行的数仓工程骨架、Kafka/HBase/ClickHouse等组件联调思路以及离线与实时链路整合的实践参考适合具备一定Spark/Flink基础、希望落地完整数仓项目的初中级开发者。1. 拿到离线实时数仓项目包先搞懂这套源码到底解决什么问题一个标着「Spark离线数仓Flink实时数仓项目源码部署资料」的压缩包解压之后大概率是两套互相独立、却共用同一套数仓分层的工程Spark 这一侧用批量任务把 ODS 数据逐层加工成 DWS 和 ADSFlink 那一侧用流式任务把同一份业务数据实时写入大宽表。我见过不少拿到这类包的人第一步就冲进 pom.xml 和 application.yml结果卡在环境上三四天跑不起来。这套项目真正值钱的地方其实不在源码而在数据模型怎么分层、两套引擎怎么对齐指标口径、以及部署参数到底怎么传。适合谁适合想把离线数仓从 Hive 迁到 Spark、或者正要补实时链路、又不想从零造轮子的数据开发。下面按「架构 → 离线 → 实时 → 排错 → 数据校验」的顺序把能直接复用的部分拆开讲。2. 离线与实时双轨架构Spark 和 Flink 在一套数仓里怎么分工2.1 批与流是两套生命周期不是两套代码离线数仓和实时数仓最大的区别不在引擎而在数据的生命周期。离线链路里数据以分区为粒度批量到达ODS 层一般按天或按小时挂新分区任务跑完就固定下来很少回头改实时链路里数据是无限流Flink 任务一旦启动就要一直挂着状态、水位线、checkpoint 全都围绕「持续运行不丢不重」设计。所以这套项目里Spark 和 Flink 被放在一起不是因为谁替代谁而是因为离线数仓负责历史回溯和全量重算实时数仓负责分钟级延迟的指标输出。两者共用同一套 ODS 数据源但从 ODS 往下就开始分叉离线侧用 Spark SQL 做大规模批量关联DWS 层的结果写回 Hive 表实时侧用 Flink SQL 做流式关联和窗口聚合结果落到 ClickHouse 这类 OLAP 引擎里供 BI 查询。很多新手拿到这种双轨项目非要把两条链路的表结构和调度周期强行对齐结果两边口径打架对不上账。2.2 ODS/DWD/DWS/ADS 四层怎么在两套引擎里对齐数仓分层的核心价值是口径统一而双轨项目最难的地方就是让离线口径和实时口径落在同一套分层语义上。我一般会先看项目里的 SQL 脚本目录确认四层在 Spark 和 Flink 两侧的命名是否成对出现。分层离线侧Spark SQL / Hive实时侧Flink SQL数据形态ODSods_xxx按天分区存原始日志或 binlog 镜像同一张 Kafka Topic作为 Source 表原样接入不加工DWDdwd_xxx 事实表 维度表宽表清洗和标准化流式 join 后的宽表可回写 Kafka明细层口径最重DWSdws_xxx 按主题聚合按天/小时累计窗口聚合结果分钟级指标汇总层指标口径统一ADSads_xxx 应用层服务报表ClickHouse 明细大表 物化视图面向查询可冗余这个表格是理解这类项目源码的总钥匙。离线 DWD 和实时 DWD 的字段定义必须同源不然流批两套指标在 ADS 层永远对不齐。常见的做法是让两份 DDL 共用同一份元数据模板离线建表语句和 Flink 的 WITH 子句都从同一张字段映射表生成。这套项目里如果离线表和实时表的字段顺序不一致跑数仓对账脚本的时候就会翻车。2.3 源码包和部署资料通常是怎么组织的我不建议拿到压缩包就到处点开看先按目录把内容归类。这类「源码部署资料」的包常见的组织方式分四块源码工程、SQL 脚本、部署文档、配置文件模板。源码工程里一般是 Spark 的 Scala 或 Java 工程、Flink 的 Java 工程还有可能带一个 SpringBoot 的调度或元数据管理小服务SQL 脚本按 ODS、DWD、DWS、ADS 四个目录分层放部署资料里通常是环境准备说明、JAR 包清单、Hive 或 YARN 的参数模板。看部署资料时我最先看三样东西Spark 和 Flink 的版本、Hive 的版本、以及 JDK 版本。这三个版本不匹配驱动冲突和序列化异常会接踵而来代码本身反而不容易出问题。如果部署文档里没有明确写版本就进 pom.xml 看 parent 和 dependencyManagement先确认引擎版本再去配置环境顺序不能反。2.4 技术选型为什么离线是 Spark实时是 Flink离线侧选 Spark 是性价比问题。Spark SQL 在 Hive on MR 的时代把查询性能提升了数倍而且 DataFrame 的 Catalyst 优化器能自动处理谓词下推和列裁剪不需要人工去调 MapReduce 的怪癖。实时侧选 Flink 是因为 Flink 的流式计算模型是纯流式的状态管理和精确一次语义比 Spark Streaming 的微批模型更成熟在做实时大宽表和窗口聚合的时候事件时间处理和迟到数据机制都更顺手。这套项目如果离线侧还在用 Spark RDD 手写 map 和 reduce那属于没发挥 Spark 的性价比合理做法应该是顶层统一写 Spark SQL底层让 Catalyst 去优化实时侧如果是用 Flink DataStream API 手动维护状态那也要评估是不是直接用 Flink SQL 更省事。数据结构简单、指标口径稳定的场景SQL 能覆盖八成以上。3. 用 Spark 跑通离线数仓分层Hive SQL 建模与 spark-submit 部署参数3.1 从 ODS 到 ADS 的 Hive SQL 骨架照着改就能用离线数仓的核心动作是建表和插数。下面这套 SQL 是这类项目里最常见的骨架写法ODS 层原样映数据DWD 层做清洗和维度退化DWS 层按主题聚合ADS 层最终服务报表。以订单场景为例-- ODS 层原始订单表按天分区保留全量镜像 CREATE TABLE if not exists ods_order ( order_id BIGINT, user_id BIGINT, product_id BIGINT, order_amount DECIMAL(10, 2), order_status STRING, create_time TIMESTAMP, etl_time TIMESTAMP ) PARTITIONED BY (dt STRING) STORED AS PARQUET; -- DWD 层清理无效订单统一枚举值拉宽维度 CREATE TABLE if not exists dwd_order_detail ( order_id BIGINT, user_id BIGINT, product_id BIGINT, category_name STRING, order_amount DECIMAL(10, 2), order_status STRING, create_time TIMESTAMP ) PARTITIONED BY (dt STRING) STORED AS PARQUET; INSERT OVERWRITE TABLE dwd_order_detail PARTITION (dt $dt) SELECT o.order_id, o.user_id, o.product_id, p.category_name, o.order_amount, CASE WHEN o.order_status invalid THEN cancelled ELSE o.order_status END AS order_status, o.create_time FROM ods_order o JOIN dim_product p ON o.product_id p.product_id WHERE o.dt $dt AND o.order_amount 0;这段 SQL 的逻辑要点在 DWD 的 INSERT OVERWRITE按 dt 分区覆盖写入保证同一天的数据重复跑不会翻倍CASE WHEN 把枚举值统一避免后续聚合时同一状态出现两种写法JOIN 维表把 product_id 拉宽成 category_name这一步叫维度退化是 DWD 层最常见的操作。跑数仓数据清洗的活儿八成都在这个环节。3.2 spark-submit 的参数模板照着传就能跑SQL 写对了只是第一步部署参数传错才是黑匣子。下面这套 spark-submit 参数模板是跑离线数仓批任务最常用的配置模板兼顾资源利用率和稳定性spark-submit \ --master yarn \ --deploy-mode cluster \ --name spark_offline_dwd \ --queue root.etl \ --driver-memory 4g \ --num-executors 8 \ --executor-memory 8g \ --executor-cores 4 \ --conf spark.sql.shuffle.partitions80 \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.sql.adaptive.coalescePartitions.enabledtrue \ --conf spark.dynamicAllocation.enabledfalse \ --class com.example.offline.DwdOrderJob \ offline-etl.jar参数最需要盯的是 queue 和 dynamicAllocation。queue 指定 YARN 队列如果项目部署文档里没建 root.etl 这个队列任务会一直挂在 ACCEPTED 状态日志里只有一串 resource manager 的提示dynamicAllocation 在生产上建议关掉开着的后果是任务跑一半 executor 被缩容shuffle 数据拉取的时候频繁报 fetch failure。shuffle.partitions 先按 10 倍 executor 核数给跑通了再根据数据量微调。3.3 Spark 内存和 shuffle 的两个必调参数Spark 离线任务的内存参数经常被误用为「调大 executor-memory 就万事大吉」。实际上 Spark 把 executor 内存分成执行内存和存储内存shuffle 数据落盘和溢写跟 spark.sql.autoBroadcastJoinThreshold 以及 spark.sql.shuffle.partitions 的关系更大。常见做法是保持默认的 spark.memory.fraction0.6但把 spark.sql.adaptive.coalescePartitions.initialPartitionNum 设成和 shuffle.partitions 一致让 AQE 在运行期把过小分区合并掉。一个实际调优顺序先看 Spark UI 的 Shuffle Read 指标如果某个 stage 的 shuffle read 大量落盘优先调大 shuffle.partitions 而不是加内存如果某个 stage 有严重的磁盘溢写再给 executor 加内存。数据倾斜则是另一套打法常见方案是把大 key 加随机前缀再二次聚合这个在项目源码里一般会体现在某个自定义 UDF 里建议先搜「skew」或「random prefix」这两个关键词。3.4 别把 Spark SQL 当 Hive 跑两个关键差异我见过最典型的翻车是把 Spark SQL 完全当 Hive 用写入时依赖动态分区插入的严格模式。Hive 对动态分区默认是 non-strictSpark 也支持但 dynamic.partitioning.enabled 默认开、pruning.enabled 默认也开很容易出现建表时没写静态分区、插入时又漏了分区字段的情况导致全表扫描。另一个差异在时间类型。Hive 的 TIMESTAMP 在 Spark 3.x 里默认走 UTC 时区如果集群时区设置是 Asia/Shanghai跑出来的时间字段会整体偏 8 小时。检查这类 Spark 数据清洗任务时我一般先看 driver 和 executor 的 user.timezone 是不是一致再用一条 SQL 直接验证select current_timestamp()和from_utc_timestamp(current_timestamp(), Asia/Shanghai)是否相等。4. Flink 实时数仓链路从 MySQL 同步到 ClickHouse 的 JDBC 连接与状态管理4.1 用 Flink SQL 把 MySQL 同步到 ClickHouse连接器与最小链路实时数仓里最常见的落地场景是用 Flink CDC 把业务 MySQL 的 binlog 同步到 Kafka再由 Flink SQL 消费 Kafka 写入 ClickHouse。这套「使用flink实现mysql同步到clickhouse」的链路在项目源码里通常以 Flink SQL 任务的方式呈现-- SourceMySQL CDC 表 CREATE TABLE mysql_orders ( order_id BIGINT, user_id BIGINT, product_id BIGINT, order_amount DECIMAL(10, 2), create_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname localhost, port 3306, username flink_user, password flink_pass, database-name business_db, table-name orders, scan.startup.mode initial ); -- SinkClickHouse 表 CREATE TABLE ch_orders ( order_id BIGINT, user_id BIGINT, product_id BIGINT, order_amount DECIMAL(10, 2), create_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:clickhouse://localhost:8123/analytics_db, table-name orders, username default, password , sink.buffer-flush.max-rows 1000, sink.buffer-flush.interval 2s ); INSERT INTO ch_orders SELECT * FROM mysql_orders;这条链路的关键在两个配置。mysql-cdc 的 scan.startup.mode 设成 initial任务启动时会先做一次全量快照再接着读 binlog增量这个模式适合中小数据量如果业务表已经几千万行initial 会把任务卡住半小时以上这种场景要改成 latest-offset 并配合手工补数。jdbc sink 的 buffer-flush.max-rows 和 interval 必须同时设置缺了任何一个ClickHouse 侧要么攒批不落盘、要么每条都发一次请求导致连接被打满。4.2 JDBC 连接器异常这 4 个参数值得优先排查搜热词里能看到「flink的jdbc连接器异常」是高频问题我在项目部署资料里也常看到对应的排查清单。JDBC 连接器最容易出问题的不是 SQL 语法而是连接参数和并发度参数建议值失效时的现象sink.buffer-flush.max-rows500-2000ClickHouse 侧数据延迟到达Task 内存涨sink.buffer-flush.interval1s-5s吞吐上不去QPS 被拉低sink.max-retries3-5网络抖动时任务直接失败重启parallelism和 ClickHouse 分片数对齐部分 shard 热点写入倾斜JDBC 连接异常还有一个隐蔽坑Flink 的 JDBC sink 默认使用单一连接parallelism 大于 1 时会创建多个连接ClickHouse 侧如果 max_open_conns 没有调大连接数会被打满。排查思路是先看 Flink UI 里对应 sink 的背压指标再看 ClickHouse 的 system.metrics 里 TCP 连接数不要在没看指标之前就盲目加 parallelism那会火上浇油。4.3 Checkpoint 与状态后端精确一次语义的两根柱子实时数仓任务里状态是流式计算的黑匣子而 checkpoint 是给状态做的后悔药。部署资料里如果给了 flink-conf.yaml我会先检查下面这段配置state.backend: rocksdb state.checkpoints.dir: hdfs://nameservice/flink-checkpoints state.savepoints.dir: hdfs://nameservice/flink-savepoints execution.checkpointing.interval: 60s execution.checkpointing.mode: EXACTLY_ONCE execution.checkpointing.timeout: 10min execution.checkpointing.min-pause: 30s execution.checkpointing.tolerable-failed-checkpoints: 3这组参数里最重要的不是 interval而是 tolerable-failed-checkpoints。生产环境网络抖动是常态如果这个值设成 0一次 checkpoint 失败就会导致任务重启重启后又要重新追数据反而更容易丢数。设成 3 的意思是允许连续 3 次 checkpoint 失败不重启给 Flink 一个自愈窗口。state.backend 选 rocksdb 是因为状态量大时堆内存装不下但 rocksdb 的代价是 CPU 序列化开销任务吞吐明显下降时优先看是否有大状态未被清理而不是急着换回 heap。4.4 事件时间与迟到数据窗口聚合的必配参数实时数仓做 DWS 层指标时窗口是绕不开的。下面这段 Flink SQL 是典型的滚动窗口聚合事件时间处理和迟到数据容忍都在里面CREATE TABLE dws_order_minute ( window_start TIMESTAMP(3), window_end TIMESTAMP(3), product_id BIGINT, order_amount DECIMAL(10, 2) ) WITH ( connector kafka, topic dws_order_minute, properties.bootstrap.servers localhost:9092, format json ); INSERT INTO dws_order_minute SELECT TUMBLE_START(create_time, INTERVAL 1 MINUTE) AS window_start, TUMBLE_END(create_time, INTERVAL 1 MINUTE) AS window_end, product_id, SUM(order_amount) AS order_amount FROM mysql_orders GROUP BY TUMBLE(create_time, INTERVAL 1 MINUTE), product_id;事件时间的核心是水位线。上面这段 SQL 没有显式设置 watermark如果源表 create_time 是 TIMESTAMP(3)Flink 会把它当成处理时间而不是事件时间窗口聚合的结果会跟实际业务时间对不上。正确做法是在 Source 表 DDL 里加 watermark 定义比如WATERMARK FOR create_time AS create_time - INTERVAL 10 SECOND给迟到的数据留 10 秒容差。4.5 Flink on YARN 的部署姿势部署资料里最常见的是 standalone 集群模式但生产上更可靠的是 Flink on YARN。区别在于standalone 模式的 JobManager 挂了要人工拉起on YARN 模式由 YARN 的 ApplicationMaster 负责拉起资源也由 YARN 统一管理。部署命令一般是两条# 启动会话集群 flink run -d -m yarn-cluster \ -yjm 1024m -ytm 4096m \ -p 4 \ -c com.example.realtime.OrderSyncJob \ realtime-etl.jar # 任务级资源隔离 flink run -d -t yarn-per-job \ -D yarn.application.nameorder_sync \ -D jobmanager.memory.process.size1024m \ -D taskmanager.memory.process.size4096m \ -c com.example.realtime.OrderSyncJob \ realtime-etl.jar区别在 -t yarn-per-job 是为这一个任务单独起一个集群任务结束集群自动释放适合实时链路里延迟敏感的核心作业。如果项目部署资料里给的命令用的是 -m yarn-cluster注意这是会话模式多个任务共用一个集群一个任务的内存溢出会把同一个集群里的邻居也带崩。选哪种取决于实时任务数量链路少用 per-job 更稳链路多用会话模式省资源但要有资源隔离的保障。5. 部署资料里常见的 5 个翻车点现象、原因与解决5.1 Flink 任务启动即报 JDBC 连接器 ClassNotFound现象任务提交到 YARN 后运行几秒钟就失败日志里出现ClassNotFoundException: org.apache.flink.connector.jdbc.JdbcSink或类似的关键字。原因JDBC 连接器的 JAR 包没有打进 fat-jar或者打进了但和 Flink 自带的连接器版本冲突。Flink SQL 客户端会懒加载连接器提交时明明没报错运行到初始化 sink 阶段才加载类这时候失败已经晚了。解决在打包插件里显式引入 flink-connector-jdbc并确认 scope 不是 provided。部署资料里的 JAR 包清单如果列了依赖对比一下版本号Flink 1.14 和 Flink 1.17 的 JDBC 连接器包名不同混用会直接 NoSuchMethodError。提交前用jar tf realtime-etl.jar | grep Jdbc验证包是否存在。5.2 Spark 离线任务跑完目标分区里的数据翻倍现象ODS 或 DWD 表的数据量每天对不上同一条订单在 Hive 里出现多条数仓对账发现离线表行数比业务库多。原因写入用了 INSERT INTO 而不是 INSERT OVERWRITE。离线任务一般会重跑重跑时 INSERT INTO 会在原有分区里再追一遍数据导致重复。尤其用 Spark SQL 提交时很多人习惯照搬 Hive 的 INSERT INTO 写法。解决把离线链路的写入统一改成 INSERT OVERWRITE配合固定分区字段。如果任务里确实需要增量追加必须在表上设计唯一的 etl_time 标记并在重跑前先 DELETE 对应分区的数据。更省事的做法是把 DWD 表建成 parquet 格式加分区覆盖反正离线数仓不追求行级更新。5.3 ClickHouse 同步任务延迟越来越高背压一直满现象Flink UI 上 sink 任务背压持续 100%Kafka 消费 lag 上涨ClickHouse 侧的写入 QPS 却不高。原因ClickHouse 是大批量写入引擎单条插入效率极低。Flink JDBC sink 如果 buffer-flush.max-rows 没设置或者设得过大数据要攒很久才 flush 一次缓冲区一直被占满背压自然下不去。解决把 flush 参数调成max-rows1000、interval2s让每个并发连接每秒写一批数据吞吐能上来一个量级。如果 ClickHouse 侧还是扛不住再看是不是 partition 字段没设置导致每个 shard 都要全量落盘把写入表的 PARTITION BY 字段和 Flink 数据分布字段对齐。5.4 小时级分区任务偶发失败日志里只有 ExecutorLostFailure现象离线任务一周内挂两三次YARN 日志里只有ExecutorLostFailure (executor lost)没有具体异常堆栈任务重启后又能跑通。原因大部分是 executor 所在节点被 YARN 回收、或者节点本地磁盘被 shuffle 中间文件写满。偶发失败且重启能成功大概率不是代码问题是资源问题。少数情况是 driver 和 executor 的 JVM 配置不一致导致 OOM。解决先看 Spark UI 里失败 executor 的 stdout/stderr 日志重点看Container exited with a non-zero exit code前后的提示。如果是磁盘问题在 spark-submit 里加--conf spark.local.dir/data1/spark,/data2/spark把 shuffle 分散到多块盘并监控节点磁盘使用率。如果是 OOM给 executor 加 memoryOverhead而不是盲目调大 executor-memory。5.5 实时任务重启后从 checkpoint 恢复不了一直跑旧状态现象Flink 任务因为发布新代码重启从 savepoint 恢复时提示状态不匹配或者恢复了但指标数据从重启前的某个时刻开始算窗口结果对不上。原因常见原因是修改了算子结构或 UID。Flink 通过 operator UID 来对应状态如果代码里没显式指定 uid重启后算子自动生成的 hash 值发生变化状态就找不到归属。第二个原因是恢复了旧的 checkpoint但 Kafka 的 offset 没有同步对齐。解决所有有状态算子都显式加.uid(order-window-agg)之类的标记从开发第一天就养成这个习惯。恢复时用 savepoint 而不是 checkpointsavepoint 是给重启准备的checkpoint 是给故障恢复准备的。保证这次重启的操作顺序停任务 → 触发 savepoint → 发布新代码 → 指定 savepoint 路径启动。6. 把两套链路拼成一条数据质量校验对账脚本与血缘回溯跑通离线链路和实时链路之后最值得花时间的不是加新指标而是把两边 DWS 层的数据做一次对账。这里的主意是让离线 30 分钟出一次结果实时 1 分钟出一次结果对账时把实时链路的结果按分钟汇聚到小时再和离线小时分区做差值比对差值超过阈值就报警。下面是这套对账脚本的 SQL 形态-- 离线侧按小时聚合订单金额 SELECT dt, hour, product_id, SUM(order_amount) AS offline_amount FROM dws_order_hourly WHERE dt 2025-01-01 GROUP BY dt, hour, product_id; -- 实时侧按小时聚合订单金额 SELECT date_format(window_start, yyyy-MM-dd) AS dt, hour(window_start) AS hour, product_id, SUM(order_amount) AS realtime_amount FROM dws_order_minute GROUP BY date_format(window_start, yyyy-MM-dd), hour(window_start), product_id;两张结果表按 dt、hour、product_id 关联算出差值率和差量。差值率超过 1% 的分钟段优先查实时链路有哪些订单没进窗口、哪些订单迟到了被水印丢弃差值长时间大于 0优先查离线链路是不是有定时任务失败导致分区缺失没补数。这套做法能同时把离线数仓的分区监控和实时数仓的延迟监控覆盖到。对账脚本做顺之后还值得做一步血缘回溯从 ADS 层的一张报表倒推确认它依赖的 DWS 表是由离线任务产出的还是由 Flink 任务产出的再往前推到 DWD 是同一张源表。用数据质量校验的差值结果反查血缘链路是最快定位业务指标对不齐的方式。我自己的习惯是每次改完数仓任务除了跑通还会手动做一次「从源表到报表」的 end-to-end 数据量核对形成习惯之后数仓的线上事故能少一半以上。希望帮到你。本文还有配套的精品资源点击获取