离线Spark与实时Flink双链路数仓架构选型与避坑指南

发布时间:2026/10/3 20:49:51
离线Spark与实时Flink双链路数仓架构选型与避坑指南 简介一套涵盖 Spark 离线数仓与 Flink 实时数仓的完整工程实现资料面向大数据开发者和数仓方向学习者适合具备常用组件基础、希望系统掌握数仓分层设计与项目落地的中高级读者。实时数仓按 ODS→DWD→DWS 分层设计ODS/DWD 层采用 Kafka 满足实时读取与分组累加DIM 层选择 HBase 做维表永久存储与按主键快速查询DWS 层使用 ClickHouse 承载汇总报表并对比 Redis、ES、Hive 等方案的选型原因。压缩包共 607 个文件、约 54.21MB以 256 个 Java 源文件为骨架配合 SQL 脚本、Shell 部署脚本、XML/JSON 配置、Markdown 说明文档及依赖配置便于按完整目录导入与调试。资源中包含订单、优惠券、购物车、用户等业务实体与服务的实现可看到业务代码如何对接 Kafka、HBase、ClickHouse部署资料能帮助快速搭建环境、复现离线与实时数仓链路。目前已有 405 人学习下载是一份可用于工程实战与面试准备的综合参考。1. 一套源码拆出两条数仓链路离线Spark与实时Flink的配合方式做数仓的人常有一种错觉源码才是最值钱的。真把一套Spark离线数仓加Flink实时数仓的压缩包拆开最该研究的反而是架构选型和分层设计——这决定了代码能不能跑得动、跑得稳。这个包覆盖了订单、优惠券、购物车、用户、退款和埋点日志几个电商核心域离线链路用Spark做传统数仓分层实时链路用Flink配合Kafka、HBase、ClickHouse完成实时加工压缩包里还带了整套部署资料从集群搭建到启动脚本都有。适合两类人一类是搭完单机Hadoop、想看看真实数仓每一层怎么衔接的入门者能照着部署资料把环境拉起来另一类是做过离线数仓、正要补实时链路的工程师能直接抄选型结论和排坑思路。2. 实时数仓选型复盘为什么是Kafka、HBase、ClickHouse而不是Redis和ES2.1 实时数仓为什么先定存储再定计算接手实时数仓项目时很多人的第一反应是选计算框架Flink还是Spark Streaming。但真正决定这条链路能走多远的往往是存储选型。计算框架只解决“怎么算”的问题存储要解决“数据放哪、怎么查、能回溯多久”的问题后者才决定了每一层能承接什么场景。这套项目里的实时链路分四层ODS用KafkaDIM用HBaseDWD继续用KafkaDWS用ClickHouse。这个分层顺序不是随便定的每一层都对应一个具体的使用场景。ODS层的场景是“每过来一条数据读取到并加工处理”要求消息队列能实时读取、实时写入DIM层的场景是“事实表根据主键获取一行维表数据”要求永久存储、主键查询DWD层是“每过来一条数据读取到并分组累加处理”本质还是流式加工所以消息队列继续扛DWS层则是面向查询分析需要列式存储支撑聚合计算。先把这四个场景列清楚再回头看选型理由就顺了。2.2 DIM层选型五个候选人逐个淘汰DIM层在实时数仓里是最容易翻车的。事实表来了一条订单数据要拿userId去维表里查出用户所在的城市、会员等级、年龄段这些信息不能每次查都全量扫描也不能靠内存硬扛。这套项目把候选人挨个过了一遍结论很干脆。Kafka被否掉的理由是不能长期存储有一些重要的用户信息需要长期保存消息队列做不到同时它不提供按主键的查询能力就算消息还在也没法直接get。Redis被否掉的理由是内存数据库用户表数据量大全量放内存的成本吃不消Hive被否掉是因为底层是HDFS查询效率低下一次维表join要等好几秒实时链路等不起ES看着能用但默认会给所有字段创建索引用户维表几十个字段全建索引写入放大严重MySQL本身压力太大真要勉强用只能挂从库减轻主库压力但量级大了还是扛不住。最后留下的HBase是典型的海量数据永久存储、按主键快速查询刚好命中“事实表根据主键获取一行维表数据”这个场景。HBase在DIM层有一个细节值得学习维表的主键就是rowkey。用户维表的rowkey直接拼userId别加随机前缀——DIM层的查询模式是点查不是范围扫描加盐反而破坏主键查询。列族设计也不要照搬业务表把经常一起查的字段放一个列族比如用户基础信息和会员等级放一起能省掉跨列族查询的开销。2.3 DWS层用ClickHouse的代价并发不是它的菜DWS层选ClickHouse在实时数仓场景里是合理选择但要注意它的边界。ClickHouse是列式存储聚合计算极快一条group by SQL能在毫秒级返回结果特别适合DWS层做轻度聚合后的多维查询。但它的并发能力很弱如果前端报表几十个用户同时点刷新直接把ClickHouse打挂是常事。这套项目把ClickHouse放在DWS层而不是直接对应用层暴露就是在规避这个短板。常见做法是在ClickHouse前面加一层查询服务把报表请求做合并和限流或者直接让BI工具连它。还有人习惯拿Redis扛用户维表数据量大的问题从这套选型来看思路反了用户表的数据量不在“查得够不够快”而在“存不存得下”内存数据库天然不适合这种场景。把存储选型的边界条件先立住再动手写Flink作业后面的麻烦会少很多。候选存储DIM层核心诉求是否满足落选原因Kafka不满足不能长期存储、不提供主键查询Redis不满足内存数据库用户表数据量太大存不下Hive不满足HDFS查询效率低下实时链路等不起ES不满足默认全字段建索引写入放大严重MySQL勉强压力太大顶多挂个从库做兜底HBase完全满足永久存储、主键高效查询3. 拆Flink实时链路从ODS到DWS的作业结构与部署参数3.1 从class名反推业务边界订单、优惠券、退款与埋点日志压缩包里那一串class名其实暴露了整套系统的业务边界OrderInfo和OrderInfoServiceImpl是订单主体与订单服务OrderRefundInfoServiceImpl是退款服务CouponInfo和CouponUseServiceImpl是优惠券及核销逻辑CartInfo是购物车UserInfo是用户维表AppAction、AppPage、AppCommon是App端埋点的行为日志和页面日志。这个模块划分在实时链路里直接映射成Kafka的topic集合。订单、退款、加购这类数据是事实数据每来一条就要加工进ODS之后直接往下游走用户、优惠券这类数据偏维表进了ODS之后要落到HBase供后续join使用AppAction和AppPage是用户行为日志是后续算转化率、漏斗的数据来源。我拿到这类包的第一步不是急着看Flink SQL而是先列一张“业务模块→数据来源→目标存储”的映射表把每个class对应的数据流搞清楚。这张表比代码更能说明问题也方便后面部署时对照topic名称。3.2 Flink作业的部署骨架并行度、Checkpoint与状态后端实时链路的代码结构可以后面慢慢看部署参数必须提前定。Flink作业能不能稳定跑很大程度取决于并行度、Checkpoint间隔和状态后端这三个参数怎么配。这套项目同时涉及Kafka读写和HBase维表查询状态后端的选择尤其关键。常见做法是生产环境用RocksDB作为状态后端因为数据量一大堆内存状态后端很快就把TaskManager的内存打满。Checkpoint间隔按业务容忍度来定允许丢几秒数据就设60秒要求严格就设30秒但别低于10秒——间隔太短Checkpoint频繁做快照反而影响吞吐。# flink-conf.yaml 关键参数 jobmanager.memory.process.size: 2048m taskmanager.memory.process.size: 4096m taskmanager.numberOfTaskSlots: 4 # 状态后端与Checkpoint state.backend: rocksdb state.checkpoint-storage: filesystem state.checkpoints.dir: hdfs://nameservice/flink-checkpoints execution.checkpointing.interval: 60s execution.checkpointing.timeout: 30s execution.checkpointing.min-pause: 40s execution.checkpointing.mode: EXACTLY_ONCE # 重启策略 restart-strategy: failure-rate restart-strategy.failure-rate.max-failures-per-interval: 3 restart-strategy.failure-rate.failure-rate-interval: 10min restart-strategy.failure-rate.delay: 10s这里重点解释三个参数execution.checkpointing.min-pause设置40秒是保证两个Checkpoint之间至少隔40秒防止连续做快照把正常处理拖垮restart-strategy用failure-rate而不是固定重启次数是考虑到实时作业偶尔会有瞬时故障3次以内自动拉起超过3次进入FAILED状态等人工介入避免无限重启掩盖真正的问题taskmanager.numberOfTaskSlots设4是给HBase维表查询留出足够的线程资源因为维表查询是典型的阻塞操作slot太少容易把作业拖成背压。3.3 维表Join HBase主键查询与缓存参数实时链路里最影响性能的往往是维表join。订单事实流里的userId要关联出用户信息如果每条数据都实时查一次HBaseRegionServer的压力会非常大。这套项目里的UserInfo和CouponInfo都是标准维表join时应该加缓存。// 维表Join HBase的常见实现片段 public class HBaseAsyncLookupFunction extends RichAsyncFunctionString, UserInfo { private Connection connection; private Table table; private CacheString, UserInfo cache; Override public void open(Configuration parameters) throws Exception { org.apache.hadoop.conf.Configuration hbaseConf HBaseConfiguration.create(); hbaseConf.set(hbase.zookeeper.quorum, zk1:2181,zk2:2181,zk3:2181); connection ConnectionFactory.createConnection(hbaseConf); table connection.getTable(TableName.valueOf(dim:user_info)); // 10000条容量、10分钟过期热点用户基本都能命中 cache CacheBuilder.newBuilder() .maximumSize(10000) .expireAfterWrite(10, TimeUnit.MINUTES) .build(); } Override public void asyncInvoke(String userId, ResultFutureUserInfo resultFuture) { UserInfo cached cache.getIfPresent(userId); if (cached ! null) { resultFuture.complete(Collections.singleton(cached)); return; } // 缓存未命中异步查询HBase查到后回填缓存 Get get new Get(Bytes.toBytes(userId)); table.get(get); // 实际生产用AsyncTable避免阻塞 // ... 解析结果并写入cache } }这段代码有三个关键点连接对象用ConnectionFactory.createConnection创建一次放在open方法里复用不能在每条数据里新建连接——那是维表查询性能杀手缓存设置10000条容量、10分钟过期这意味着热点用户的维表信息10分钟内不会重复查HBase只有首次查询或缓存过期后才走网络IO查询用AsyncTable而非普通Table因为Flink的异步IO要求不能阻塞等待结果同步查询会卡住整个算子。缓存的容量和时间要根据维表数据更新频率来调用户信息一天更新一次缓存设大点没问题优惠券信息如果运营会临时改状态缓存过期时间就要缩短到几分钟否则改券状态后实时计算里还是旧值排错时非常容易误判成代码问题。4. 拆Spark离线链路分层建模与调度提交4.1 离线数仓的每一层职责ODS、DWD、DWS、ADS怎么切离线链路用Spark而不用Hive核心原因是Spark SQL的执行速度更快同样跑一张大表关联Spark比MapReduce快一个量级。分层架构还是老四样ODS、DWD、DWS、ADS这套项目里的订单、用户、优惠券、退款、埋点日志在离线链路里能找到对应的影子只是处理逻辑和实时链路截然不同。离线ODS层是原样落地业务库的MySQL binlog和App端埋点日志落成Hive表保留原始数据做最简单的分区和压缩DWD层才是清洗和建模的主战场订单表要拉平、去重、补齐维度字段用户表要做拉链表处理历史变化DWS层做轻度聚合按用户、商品、日期几个常用维度汇总结果供下游报表直接查询ADS层是应用层针对具体业务需求加工比如大屏指标、运营日报。这个分层和实时链路的分层可以一一对应实时ODS对应Kafka离线ODS对应Hive表实时DIM对应HBase离线DWD里的维表就是普通Hive维表。一个关键的认知是离线数仓处理的是T1全量数据实时数仓处理的是秒级增量数据。同一张订单表离线DWD要去重昨天的全量数据实时DWD是每来一条加工一条。所以离线DWS里的订单金额汇总和实时DWS里的订单金额汇总天然存在一个“时间窗口覆盖差异”这一点在第五章展开讲是最隐蔽的对不上账的根源。4.2 用Spark SQL建分层表订单事实与用户维度的DDL实践离线链路里建表是最容易草率、也最容易在后期吃大亏的环节。字段类型没选对、分区设计不合理跑数的时候才会暴露。以订单事实表和用户维表为例这套项目的DDL里值得注意几个细节。订单事实表用日期作为分区字段按天分区数据量大的时候再叠加hash分桶金额字段用DECIMAL(14,2)而不用DOUBLE避免浮点误差导致对不上账——这是离线数仓的必修课就不多说了。用户维表需要处理历史变化用拉链表设计start_date和end_date标识有效期。-- 订单事实表 DWD层 CREATE TABLE IF NOT EXISTS dwd_order_info ( order_id BIGINT COMMENT 订单ID, user_id BIGINT COMMENT 用户ID, coupon_id BIGINT COMMENT 优惠券ID, order_amount DECIMAL(14,2) COMMENT 订单金额, pay_amount DECIMAL(14,2) COMMENT 实付金额, order_status TINYINT COMMENT 订单状态, province_id BIGINT COMMENT 省份ID, create_time TIMESTAMP COMMENT 下单时间 ) PARTITIONED BY (dt STRING COMMENT 日期分区) STORED AS PARQUET; -- 用户维表 拉链表 CREATE TABLE IF NOT EXISTS dim_user_info ( user_id BIGINT COMMENT 用户ID, user_name STRING COMMENT 用户姓名, phone_num STRING COMMENT 手机号, member_level TINYINT COMMENT 会员等级, register_time TIMESTAMP COMMENT 注册时间, start_date STRING COMMENT 生效日期, end_date STRING COMMENT 失效日期9999-12-31表示当前有效 ) STORED AS PARQUET;第一个建表语句里order_status用TINYINT而不是STRING是为了将枚举字段沿用业务库的编码方式DWD层先不做翻译翻译动作可以放到DWS或ADS层。第二个建表语句是拉链表的标准写法离线每日任务跑完后把当天有变更的用户闭链、新用户开链查询时用WHERE dt2025-01-15 AND user_id123就能拿到用户在那天的实时画像。这类DDL在部署资料里一般都会带但直接复制容易忽略一个问题拉链表如果不配合每天的更新任务只建表不跑数据那这个维度表跟普通全量表就没有区别。4.3 调度与提交一个能跑的Spark任务长什么样离线链路跑不跑得起来多半看调度和提交参数。Spark作业不能靠人肉点击submit要用调度工具按天触发。最常见的做法是crontab配合shell脚本凌晨一点跑昨天的数据代码写成shell提交Spark SQL任务失败重试一次重试还失败就发告警。#!/bin/bash # spark离线任务每日调度脚本 TODAY$(date %F) YESTERDAY$(date -d yesterday %F) # 先检查ODS分区是否就绪 hadoop fs -test -e /warehouse/ods_order_info/dt$YESTERDAY if [ $? -ne 0 ]; then echo ODS分区不存在任务终止: $YESTERDAY exit 1 fi spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 8g \ --num-executors 20 \ --conf spark.sql.shuffle.partitions200 \ --conf spark.dynamicAllocation.enabledfalse \ --class com.example.dwd.DwdOrderJob \ /data/jars/offline-etl.jar \ --dt $YESTERDAY这个脚本里有三个惯用做法值得说明。hadoop fs -test -e先检查上游ODS分区是否存在避免上游数据没就绪时下游任务白跑一趟这是数仓调度里最简单也最实用的防御手段。spark.sql.shuffle.partitions200是经验值分区数太少会出现数据倾斜太多则小文件爆炸200个分区对单表几千万到几亿的数据量是一个合理起点。spark.dynamicAllocation.enabledfalse是固定资源避免动态分配在凌晨调度高峰期抢不到资源导致任务排队。离线链路的坑很多不在代码里而在调度上下游的衔接。ODS表分区依赖数据同步任务DWD表分区依赖ODS分区DWS依赖DWD一环没跑下游全部白等。所以部署资料里的调度配置必须带依赖检查只要缺分区就立刻终止并告警而不是继续跑一个注定失败的作业。5. 双链路落地避坑五个真实翻车现场与解决记录5.1 维表查询变慢HBase连接没复用缓存穿透成雪崩现象实时作业刚上线时吞吐正常跑了两个小时后TaskManager持续背压HBase RegionServer的CPU打满Kafka消费延迟从秒级涨到分钟级。原因维表join的代码里在每条数据上创建了新的HBase连接连接没复用大量TIME_WAIT连接把RegionServer拖垮。还有一个隐蔽问题热key用户反复查缓存不命中就全打到HBase上缓存穿透加剧了雪崩。解决连接在open方法里创建并复用缓存用maximumSize10000加expireAfterWrite10分钟兜住热点查询HBase走Flink异步IO限制最大并发请求数给HBase留出喘息空间。从那以后我在代码评审里看到new Connection出现在算子内部一律打回重写。5.2 ClickHouse被报表并发打爆分析引擎不是服务引擎现象DWS层数据落到ClickHouse后业务方直接拿JDBC连上去做报表上线第一天下午3点ClickHouse CPU 100%查询队列堆积实时大屏数据卡住不刷新。原因ClickHouse的定位是分析引擎擅长跑大查询但并发能力弱。十几个人同时拖报表每个查询都触发全表扫描CPU当场被占满。解决在ClickHouse前面加一层查询服务做请求合并、限流和结果缓存报表场景改成预聚合查询把常用指标提前算成结果表避免每次实时聚合。这套项目的分层里DWS用ClickHouse而不让应用直接连它就是在用架构规避并发问题。5.3 Kafka消费延迟飙升并行度和分区数不匹配现象Flink作业从Kafka消费订单数据起初延迟稳定在秒级数据量上来后延迟飙到十几分钟怎么调并行度都不见好转。原因Kafka topic只有3个分区Flink作业并行度调到20但source算子的并行度最大只能是3其余19个并行度在source阶段全部闲置数据全卡在3个分区里排队。并行度不是越大越好它受上游分区数硬性约束。解决先确认topic分区数再设并行度。常见做法是topic分区数设成目标并行度的1.5倍到2倍给水平扩展留余量同时确保每个分区有独立的消费线程。部署资料里的Flink参数表如果只写了并行度没写分区数建议自己把两者对齐再上线。5.4 离线实时对不上账口径差异比代码bug更隐蔽现象离线数仓DWS的订单金额和实时数仓DWS的订单金额对不上离线比实时少了一截但两边单查都正确代码也没发现明显bug。原因离线链路按T1跑凌晨1点捞的是昨天00:00到23:59:59的数据实时链路按事件时间聚合晚到的数据会被归到实际发生的那一天。如果业务系统在23:50还在补单离线任务大概率漏掉了这批跨天数据两边口径自然不一致。解决建立统一的口径定义明确迟到数据的归属规则然后写对账脚本每天比对差异率超过阈值就告警。第六章会给出一个可用的对账脚本模板。5.5 部署资料和代码版本对不上先看class名再动手现象按部署文档搭完环境启动Flink作业时报类找不到打开jar包一看里面class名和文档里的模块列表根本对不上。原因压缩包里的class文件是编译产物可能来自多个版本的迭代部署文档写的版本比class更旧。接手这类资源时最忌讳直接照文档启动。解决先把class名列出来反推实际包含的模块再和部署文档对照确认版本匹配后再动手。这个包里的OrderInfoServiceImpl、OrderRefundInfoServiceImpl等class名已经说明了业务边界以实际产物为准以文档为辅顺序别搞反。6. 用对账脚本验证双链路离线实时口径统一的最后一公里双链路都跑通之后最值得做的一件事是写对账脚本。每天凌晨离线任务跑完拿离线DWS的结果和实时DWS的结果做一次对比差异超过阈值立刻查原因。这个脚本能同时验证两条链路的健康度也能倒逼口径统一。-- 对账SQL按日期和省份维度对比离线与实时订单金额 SELECT off.dt, off.province_id, off.order_amount AS offline_amount, real.order_amount AS realtime_amount, ROUND((off.order_amount - real.order_amount) / NULLIF(real.order_amount, 0) * 100, 2) AS diff_rate FROM ( SELECT dt, province_id, SUM(pay_amount) AS order_amount FROM dws_order_info WHERE dt ${YESTERDAY} GROUP BY dt, province_id ) off FULL OUTER JOIN ( SELECT DATE_FORMAT(ts, yyyy-MM-dd) AS dt, province_id, SUM(pay_amount) AS order_amount FROM dws_realtime_order_agg WHERE DATE_FORMAT(ts, yyyy-MM-dd) ${YESTERDAY} GROUP BY DATE_FORMAT(ts, yyyy-MM-dd), province_id ) real ON off.province_id real.province_id AND off.dt real.dt WHERE ABS(off.order_amount - real.order_amount) / NULLIF(real.order_amount, 0) 0.01;对账的阈值一般设1%超过就报警。常见差异来源就三类实时链路数据延迟导致当天数据没聚合完离线任务抽取时业务库还有未提交事务口径定义不一致——比如离线把退款金额单列实时把退款冲减到订单金额里。找到差异后不要急着改代码先定位是哪种类型按类型处理。这个包解压后我建议你也先做这四步列class名反推业务模块、对照部署资料确认版本、按第三章和第四章的参数检查作业配置、最后把一条业务线比如订单金额的对账跑通。做完这四步再往里面填业务细节比直接跑demo再返工省力得多。从那以后我每次交付数仓项目第一件事不是看报表漂不漂亮而是强制走一遍对账脚本确认两条链路的口径能对上才开始谈业务指标。希望帮到你。本文还有配套的精品资源点击获取