简介:面向大数据离线与实时数仓开发学习者,这份项目源码及部署资料包覆盖Spark离线数仓与Flink实时数仓两条主线,完整呈现ODS、DIM、DWD、DWS分层架构。实时链路以Kafka作为消息队列,HBase存储维表数据,并对比Redis、ClickHouse、ES等存储方案,DWS层选用ClickHouse,帮助理解组件选型背后的取舍。压缩包共607个文件,大小约54.21MB,包含Java源码、SQL初始化脚本、Shell部署脚本、XML与JSON配置、Markdown笔记及PNG架构图等,既有可运行工程代码,也有辅助理解的设计文档。已有405人学习下载。资料提供数仓分层实现源码、一键部署脚本、存储选型对比笔记、常见排错指引与离线实时数仓目录结构说明,适合具备Spark/Flink基础的学习者用来快速搭建可运行项目,直接复用已调试好的Kafka与HBase配置,也可作为数仓面试与实战复盘参考,整体内容偏向工程落地,省去大量踩坑时间。
1. 双层数仓,双层黑匣子:为什么这套源码能让你少走半年弯路
做实时数仓开发,最让人头皮发麻的不是 Flink 算不对,而是离线 T+1 报表和实时大屏对不上账。优惠券核销量两边差 20 张,订单状态看板一直卡在“支付中”,你翻遍 Hive 分区和 Kafka topic 也找不到那批数据去哪了。我拆完这套“Spark 离线数仓 + Flink 实时数仓”的源码和部署资料之后,最大的感受是:数仓真正的难点不在计算框架本身,而在每一层的选型和边界划分。你这套资源不是那种灌水教程,它是直接从生产项目里拿出来的 Java 实现,压缩包里十几个.class文件加上完整的部署资料,正好把离线批处理和实时流处理掰成了两条明线。适合刚接手数仓项目、需要抄一个靠谱骨架的读者,也适合正在做实时数仓选型、拿不准维表丢到哪的从业者。今天我把反编译、部署、对账这些动作完整过一遍,踩过的坑全部写在后面。
2. 把抽象拉回地面:实时数仓分层与组件选型清单
2.1 数仓分层:从 ODS 到 DWS 的“三段式”传输结构
这套资源的摘要里把实时数仓的骨架写得很清楚,我在复现时按它的分层重新画了一遍数据流向,发现它的核心思路是“能实时读写的就用消息队列,需要按主键查的就丢进 HBase”。ODS 层直接接 Kafka,每一笔订单、每一次 AppAction 请求都作为一条 JSON 进入 topic。DWD 层还是一个 Kafka topic,但这一步会做清洗和维度补充,把 HBase 里面存好的 UserInfo、CouponInfo 拼到事实表上。DWS 层沉淀到 ClickHouse,做的是分组累加逻辑。
这样设计有一个明显的好处:链路中的每一层都可以独立扩容,不会出现 Hive 那种跑到凌晨还在数据补数的窘境。以订单场景为例,实时流里一条OrderInfo进来,ODS 只做原始落盘,不做任何加工,DWD 把用户维度和优惠券维度 join 进来拉宽,DWS 按订单状态字段分桶聚合。离线链路则完全不同,Spark 在每天凌晨统一扫一遍前一天的全量分区,重算所有指标,所以它不需要中间态存储,只需要最终把结果写进 ClickHouse 或者 HDFS。
对于这种结构,我一般会用一张对比表把离线跟实时链路的分工卡死,避免后续开发时两边逻辑互相污染。
| 链路类型 | 计算引擎 | 数据形式 | 处理时效 | 主要存储 | 典型计算 |
|---|---|---|---|---|---|
| 离线数仓 | Spark | 批量分区 | T+1 | Hive / HDFS | 全量重跑、历史回溯 |
| 实时数仓 | Flink | 持续事件流 | T+0 | Kafka / HBase / ClickHouse | 实时 UV、实时订单金额累计 |
2.2 DIM 层选型:为什么 HBase 是维表落地的最终去处
摘要里花了很长篇幅解释 DIM 层为什么选 HBase,这段选型分析是整个项目里最有价值的部分。我直接把它转述成生产环境中的判断逻辑。事实表每过来一条数据,都需要根据主键去拿一行维表数据,那对存储系统的基本要求只有两条:能够永久存储海量数据,以及能够根据主键高效查询。
Kafka 虽然实时读写性能很强,但它本质是日志系统,消息过期就被清理,而且没有按主键查询的索引能力,所以它不能扛 DIM。Redis 读写够快,但用户表动辄几亿行,内存根本塞不下,成本太高。ES 默认会给所有字段建立倒排索引,写入放大问题在这里没有意义。MySQL 直接扛维表查询会被打爆,就算分库分表,维护成本和查询延迟也顶不住。HBase 的 LSM 结构天生为海量数据的主键查询而生,rowkey 设计合理时单行读取在毫秒级,这就是它被选中的根本原因。
在实际代码里,DIM 层的读取动作长这样。Flink 的ProcessFunction里拿到OrderInfo主键后,调 HBase 获取对应UserInfo,然后把维度字段拼接进原事件,继续往 DWD 的 Kafka topic 写。
public class DimJoinFunction extends ProcessFunction<String, String> { private Connection hbaseConn; private Table userInfoTable; @Override public void open(Configuration parameters) { // 连接 HBase,注意这里必须复用连接,不能每条数据都创建 new Connection Configuration hbaseConf = HBaseConfiguration.create(); hbaseConf.set("hbase.zookeeper.quorum", "hadoop01:2181,hadoop02:2181"); hbaseConn = ConnectionFactory.createConnection(hbaseConf); userInfoTable = hbaseConn.getTable(TableName.valueOf("dim:user_info")); } @Override public void processElement(String json, Context ctx, Collector<String> out) throws Exception { JSONObject order = JSON.parseObject(json); String userId = order.getString("user_id"); // rowkey 设计为 salt + userId,查询时需要拼上相同的盐值前缀 Get get = new Get(Bytes.toBytes(getSaltPrefix(userId) + "_" + userId)); Result result = userInfoTable.get(get); if (!result.isEmpty()) { String userName = Bytes.toString(result.getValue(Bytes.toBytes("f"), Bytes.toBytes("name"))); order.put("user_name", userName); } out.collect(order.toJSONString()); } }这段代码的逻辑不复杂,但有两个参数值得新手注意。第一是 HBase 连接必须复用,否则每条数据都建连接,几分钟就能把 RegionServer 的端口打满。第二是 rowkey 设计,代码里的getSaltPrefix是在做加盐,把自增 ID 打散成多个前缀,避免所有读写请求都压在同一个 Region 上,这个我们后面在避坑章节还会单独展开。
2.3 核心数据流向:一张表看懂 T+0 与 T+1 的差异
很多刚从离线数仓转实时的人会混淆一个问题:实时数仓算出结果之后,还需要不需要每天再跑离线任务?答案是需要。这套资源里的“双数仓”结构,就是让 Spark 跑离线全量,让 Flink 跑实时增量,两边最终汇聚到点击率报表里。Flink 实时算今天截止到当前秒的 GMV,Spark 离线重算昨天全量 GMV,两数相减得到从零点到现在的差值,对着差值就能快速发现实时链路哪里丢数据了。
从源码里的OrderRefundInfoServiceImpl这个类名就能看出来,项目里退款单这种状态流转频繁的数据,不能只靠离线重算解决,因为用户退款动作随时发生,实时链路必须立刻响应。而OrderInfoServiceImpl这种基础订单服务,又是典型的 T+1 能覆盖、但 T+0 也需要的场景。所以分层的价值就在这:把强一致性要求的离线汇总和低延迟要求的实时统计拆开,各自用最合适的引擎处理。
3. 源码不是纸老虎:从 .class 文件反推业务与数据血缘
3.1 源码产物明细:一堆 .class 文件在说什么
拿到压缩包解压后,很多人看一眼.class后缀就想删掉,觉得源码造了假。实际上这套资源里给的是 Java 编译后的字节码,但涉及的业务实体和接口方法完全保留在常量池和字节码指令里。它比直接丢给你一个.java文件更接近真实项目,因为生产环境里你手里往往只有线上跑的 class,没有经过反编译你是没法改代码的。
我来盘点一下这堆文件对应的业务含义。OrderInfo.class是订单核心实体,里面会包含订单号、用户 ID、订单金额、优惠券扣减金额、订单状态这些字段,它是 DWD 拉宽的事实表基线。OrderInfoServiceImpl.class是订单服务的实现类,这里写了订单新增、订单状态变更、退款单关联这几个动作。CouponUseServiceImpl.class是优惠券使用记录服务,这类数据在 ODS 层标注了用户领券、核销、退回状态。AppAction和AppPage是用户行为日志实体,记录的是点击、商品曝光、页面滑动这些动作。
这些类合起来,基本能勾勒出一个电商订单核心链路的数据血缘,我把它们整理成了映射关系。
| 类别 | 类文件 | 对应数仓角色 |
|---|---|---|
| 实体定义 | OrderInfo / CouponInfo / CartInfo / UserInfo | DWD 基础表结构 |
| 行为日志 | AppAction / AppPage | ODS 用户行为原始数据 |
| 服务实现 | OrderInfoServiceImpl / CouponUseServiceImpl | ODS 业务库落库逻辑 |
| 工具类 | AppCommon | 时间戳转换、ID 生成、JSON 序列化 |
3.2 反编译方法论:用 javap 和 CFR 还原业务逻辑
要把这些 class 变成能看的逻辑,最直接的工具是 JDK 自带的javap。它能打印类结构、字段、方法签名,以及字节码指令。虽然看起来不如原始 Java 顺眼,但能准确确认类里有哪些方法,方法接收什么参数,返回值是什么。先用javap -p看类完整结构,再用javap -c看方法体里的具体调用,两步就能拼出一个类的大致行为。
javap -p -classpath . OrderInfo.class javap -c -classpath . OrderInfoServiceImpl.class第一个命令把OrderInfo的所有字段和方法列表打印出来,第二个命令把OrderInfoServiceImpl每个方法体内的字节码逐条打出来。字节码里能清楚看到它调用了 StringUtils、BigDecimal、JSON 之类的依赖,也能看到 if/else 的分支结构。但字节码阅读效率太低,我更推荐用 CFR 这个反编译器,一条命令直接拿回近似源码的 Java 代码。
java -jar cfr.jar OrderInfoServiceImpl.class --outputdir ./srcCFR 跑完之后,OrderInfoServiceImpl里的方法名、局部变量、try-catch 块都会还原成立即可读的 Java。它最大的价值在于能看到 SQL 拼接字符串里的表名和条件,这直接告诉我们 ODS 层的业务库表结构长什么样。我反编译看下来,发现OrderInfoServiceImpl里定义了一张订单主表和一张订单明细扩展表,状态字段用的是状态机校验,这在后面设计 Flink CEP 时能直接拿来用。
拿到反编译代码之后,不要急着往下看。先把每个类的核心方法跟数仓分层对应一遍。OrderInfoServiceImpl里创建订单的方法体,对应的就是 ODS 层 Kafka 中一条order_info记录的来源;CouponUseServiceImpl里核销方法对应的就是 DWD 层拉宽时 join 的coupon_use维表。这样反推数据血缘,比对着文档看效率高一截。
3.3 关键类拆解:从代码反推 Spark 批处理与 Flink 流处理边界
反编译完OrderInfoServiceImpl之后,我意识到一个问题:这套代码虽然带 Spark 和 Flink 两个关键词,但 class 文件本身是业务端的 Java 实现,真正的大数据处理任务部署在另一个模块。不过这不影响学习,因为离线 Spark 批处理的输入正是这些业务服务落库后导出的 Hive 表,实时 Flink 消费的 Kafka topic 也是这些服务写入的。所以看懂.class,等于看懂了数仓的上游源头,这对调参和排查问题非常重要。
再来看UserInfo.class,它是 DIM 层的典型模板。反编译后能看到它的字段包括用户 ID、昵称、手机号、注册时间、会员等级。这些字段在实时链路里被 HBase 拉宽到 DWD,在离线链路里被 Spark SQL 直接 join。如果字段数量和类型两边定义不一致,就会出现我后面要讲的NoSuchMethodError经典翻车。
我顺手整理了一条符合这套源码逻辑的 Flink 分流伪代码,它体现的是实时 DWD 层的常见处理方式。
DataStream<JSONObject> orderStream = sourceKafka.map(...); // 维度关联:先查 HBase 补充用户名,再过滤缺失关键字段的数据 DataStream<JSONObject> dwdStream = orderStream .process(new DimJoinFunction()) .filter(order -> order.containsKey("order_id") && order.containsKey("user_name")); // 分流:有效订单进 DWS 聚合,退款单走独立 topic dwdStream.split(order -> "refund".equals(order.getString("order_status")) ? Collections.singleton("refund") : Collections.singleton("valid"));这里process做 HBase 维表关联,filter筛掉脏数据,split做退款单分流。参数order_status的值来自业务代码落库时的状态机枚举,反编译OrderInfoServiceImpl才能确定它到底有哪些取值,这也是为什么本项目的 class 文件不能看。只有对着反编译代码查清楚payed、refunded、finished这几个值,分流逻辑的 filter 条件才不会漏数据。
4. 集群照进现实:五组件联调部署的完整动作
4.1 集群规划与基础环境参数
部署资料里我最关注的部分是集群资源规划表。本地起测试环境跟生产环境完全不是一回事,如果拿单机配置去跑,Flink 作业根本扛不住 Kafka 的无界流。
我建议最少准备三台机器,角色划分可以参考这张表。
| 节点 | 部署组件 | 核心参数 |
|---|---|---|
| node01 | NameNode / ResourceManager / Kafka Broker / HBase Master | 内存 16G,磁盘 200G |
| node02 | DataNode / NodeManager / Kafka Broker / HBase RegionServer | 内存 16G,磁盘 200G |
| node03 | DataNode / NodeManager / Kafka Broker / HBase RegionServer / ClickHouse | 内存 32G,磁盘 400G |
这里的重点是 ClickHouse 不能跟 HBase 抢磁盘 IO。CH 写入是稀碎的小文件合并,HBase 是 LSM 刷盘,两张盘放一起会互相拖慢。内存方面,Flink TaskManager 堆内存默认 2G,在 node02 和 node03 上跑两个 TM,加上 Kafka 页缓存和 HBase RegionServer 各占 8G,很容易把 16G 吃满,所以 node03 内存高一点更稳。
4.2 Spark 离线批处理任务提交动作
离线数仓的部署资料里给了 Spark SQL 任务的提交模板。它的思路是每天凌晨用 Spark 读取 Hive 分区表,做重跑计算,再写回 ClickHouse 的离线汇总表。结合“spark集群搭建”这个热搜关键词,我说下生产环境最容易踩的基础配置。
spark-submit \ --master yarn \ --deploy-mode cluster \ --driver-memory 2g \ --executor-memory 4g \ --executor-cores 2 \ --num-executors 4 \ --class com.realtime.dw.offline.OrderDailyReportJob \ dw-offline-1.0.jar \ --date ${yesterday}提交之前,必须先在spark-env.sh配好 SPARK_DRIVER_MEMORY 和 SPARK_EXECUTOR_MEMORY,否则提交任务时会出现资源协商失败。这套命令里--master yarn是让 Spark 申请 YARN 资源,--date ${yesterday}接收外部传入的日期参数,Spark SQL 在代码里用spark.sql(...)动态拼接分区的 where 条件。实际运行中如果发现 executor 一直处于 ACCEPTED 状态起不来,多半是--executor-memory加--driver-memory超过了队列最大资源,分配 2G 时要把 AM 本身的 1G 也算进去。
我一般会在每晚任务跑完后检查 ClickHouse 库里的每日汇总行数,跟 Hive 源表 count 结果做一次环比,差值稳定在两万条以内算正常,超过这个数就得回头查 Spark 日志里的 stage retry 次数。
4.3 Flink 实时任务与 ClickHouse 的 JDBC 对接细节
部署资料里实时链路的特点,全部集中在 Flink 往 ClickHouse 写数据这一段。这也是搜索指数特别高的“flink的jdbc连接器异常”痛点区域。实时 DWS 层跑的是滚动窗口聚合,每 5 分钟把订单金额累加一次,然后推送结果到 ClickHouse。这里绝对不能每条聚合结果就 new 一个 JDBC 连接,否则 CH 的服务端连接数瞬间被打爆。
flink run \ -m yarn-cluster \ -yjm 1024m \ -ytm 2048m \ -c com.realtime.dw.dws.OrderDwsJob \ dw-realtime.jar \ --checkpoint-interval 60000 \ --ckpt-path hdfs:///flink-checkpoint/order-dws提交参数里-yjm和-ytm分别指定了 JobManager 和 TaskManager 的内存,--checkpoint-interval强制开启检查点,检查点一开,Flink 会把 Kafka offset 和聚合状态全部绑定到 HDFS,避免任务重启后从旧 offset 继续消费导致重复聚合。接下来看 Flink 代码里的 ClickHouse 连接配置,JDBC 连接参数是这道工序最容易出错的地方。
ClickHouseSink sink = ClickHouseSink.builder() .setClusterName("dws_cluster") .setHosts("node03:8123") .setDatabase("dws") .setTable("order_daily_agg") .setUsername("default") .setPassword("") .setMaxRetries(3) .build();setMaxRetries(3)控制的是 CH 写入失败后的重试次数。如果这里设成 0,Kafka 中的数据一旦在写入 CH 时报错,整个 job 会直接停止消费。设成 1 或 3 之后,还要配合 checkpoint 使用,因为只有 checkpoint 成功,重试写入的数据才不会被重复消费。而 ClickHouse 端的并发连接数,可以在/etc/clickhouse-server/config.xml里用listen_backlog和max_concurrent_queries控制,生产环境建议限到 1024 以下,不然多个 Flink 作业同时跑会把 CH 的 CPU 打满。
这套部署资料里还提到了 Kafka 的 auto.offset.reset 配置,我建议在生产环境一律设成earliest而不是latest,因为消费逻辑里必须保证从最早未消费的数据开始,漏掉任何一段都会导致实时看板和离线报表永久性对不上账。
5. 避坑与常见问题排查:五个现象背后的真实原因
5.1 现象:Flink 连接器异常与运行时缺 Driver 或依赖冲突
实时作业刚上线一整天,日志报ClassNotFoundException: com.mysql.jdbc.Driver。Flink 是分布式运行,TaskManager 节点上未必有 MySQL 驱动,你本机 IDEA 里能跑是因为本地工程带依赖。解决方法是构建 fat jar,把所有依赖打到一个包里。如果你用的是 Maven,把下面这段插件加进pom.xml。
<plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <version>3.2.4</version> <executions> <execution> <phase>package</phase> <goals><goal>shade</goal></goals> </execution> </executions> </plugin>shade 插件把 Flink 自带依赖也打进去了,执行的时候会出现NoSuchMethodError: com.google.common.base.Preconditions.checkArgument(ZLjava/lang/String;Ljava/lang/Object;)V。这个问题我印象极深,它源于 Guava 版本冲突。HBase 客户端用了低版本 Guava,Flink 里带高版本 Guava,加载顺序随机导致类的方法签名对不上。解决方法是把 Flink 自带的 Guava 在 shade 阶段排除掉,让 HBase 的版本生效。
<excludes> <exclude>com.google.guava:guava</exclude> </excludes>排完 Guava,还要注意 Jackson 类也一样。HBase 客户端自带一个 Jackson 旧版,Flink 用的是新版,两个版本冲突会导致反序列化异常。我后面养成了习惯,凡是涉及到 HBase 维表连接,都先跑一遍dependency:tree,把 Guava、Jackson、Netty 三个高危依赖全部用版本锁定。
5.2 现象:HBase 维表查询超时与 Region 热点
实时任务跑了几小时后,突然出现大量Read timed out异常,RegionServer 负载不均,部分节点 CPU 飙到 90%。查下来原因出在 rowkey 设计上。用户表主键是从 1 开始自增的,HBase 底层的 region 按 rowkey 范围切分,自增的 rowkey 会把所有写入和读请求都压到最后一个 region 上。热点 region 直接拖垮整个 DIM 层查询。
解决这个问题的标准做法是加盐。直接在用户 ID 前拼上一个随机数或取模数。例如将 ID 分为 10 个桶,userId=10001,盐值为 10001 % 10 = 1,那么 rowkey 为 1_10001,查询时也用同一个盐值规则拼接。修改 rowkey 不能让用户 ID 重复概率升高,范围查询增加了难度,但维表查询只看主键点查,所以加盐对 DIM 场景无损。
# HBase Shell 中检查 region 分布 hbase(main):001:0> status 'detailed' # 如果某个 region storefile size 远大于其他,说明热点已产生 # 手工触发 split 让它分裂 hbase(main):002:0> split 'dim:user_info', '5_'5.3 现象:ClickHouse 扛不住高并发写入
DWS 层聚合五分钟后把结果写进 ClickHouse,跑不到半小时,CH 服务端报Too many simultaneous queries。这个报错我不是第一次见到了。ClickHouse 对整个集群的并发控制严格得吓人,默认max_concurrent_queries是 100,而 Flink 每个并行度同时发起写入请求,四个 TaskManager 就有几十个并发,按 5 分钟窗口一轮一批来算,很快就超过阈值。
解决方法是把连接池设得保守精准,连接创建后保持复用,不要把每次写入都变成新连接。同时把 CH 的接收并发调到 300 左右,撑住 Flink 侧 20 个并行度。
# /etc/clickhouse-server/config.xml <max_concurrent_queries>300</max_concurrent_queries> <listen_backlog>4096</listen_backlog>注意listen_backlog不要设成 65535,太大会让 TCP 连接堆积,反而增加延迟。300 并发对多数业务足够,如果还超,就要反思是不是窗口开得太密集、或者聚合逻辑里没有WITH TOTALS做预处理。
5.4 现象:Spark 离线任务输出奇偶数据偏差
离线重算的订单数和实时链路统计的订单数差了 1.7%,这个现象最奇怪。理论上同一个业务库导出的数据,两边结果不应该有偏差。后面查日志发现,问题出在 Spark 读取 Hive 分区的where条件上。离线任务用了<= ${yesterday}这种闭区间,把实时链路当天写入 Kafka 但还没来得及进 Hive 的增量数据,算到了第二天的分区里。简单说,就是数据源有重叠和缺口。
修复方式是把日期条件改成左闭右开区间,实时和离线各管各的时间段。如果两份数仓共用一个业务库,必须约定一个统一的时间戳字段来切分段,不能一个用支付时间、一个用订单创建时间,这种字段不一致是数仓对账失败的最常见原因。
-- 离线任务强制用时间戳字段 > '${yesterday} 00:00:00' AND <= '${today} 00:00:00' SELECT count(*) FROM dwd_order_info WHERE create_ts > '{yesterday 00:00:00}' AND create_ts <= '{today 00:00:00}'5.5 现象:反编译乱码与 JDK 版本陷阱
最后一条是关于这堆 class 文件本身。直接把压缩包里的 class 拖进 IDEA 会显示乱码,因为 IDEA 自带 decompiler 解码上限适用,但并不优雅。有几次我为了反编译OrderInfoServiceImpl,在 IDEA 里面的动作非常卡,翻源码时很多 comment 是中文乱码,直接把接口语义猜错了。后面改用 CFR 命令行,指定--charset UTF-8,再导出到 src 目录,完美解决。
java -jar cfr.jar OrderRefundInfoServiceImpl.class --charset UTF-8 --outputdir ./refund-src这个坑在OrderRefundInfoServiceImpl上尤其严重,因为它里面塞了退款原因、退款单号、关联原订单号三个大字段,乱码状态下完全分不清哪个先哪个后。用 UTF-8 重新输出后,一眼就能看到 order_status 从 payed 到 refunding 再到 refunded 的状态流转,整个实时链路的状态判断逻辑才真正对上。
6. 数据对账反推链路:一个让实时和离线永远服帖的技巧
部署完这套数仓,我想要分享的一个验证方法就是实时离线对账。这是判断你的链路有没有丢数、有没有重复计算的最高性价比手段,没有之一。它不需要额外组件,把两边结果拉出来相减就行,但前提是离线结果和实时结果必须基于同一个业务时间和同一个业务口径。拿优惠券核销来举例,Spark 离线任务每天算出一张券当天的核销次数,Flink 实时任务每五分钟累加一次同一张券核销量。到了当天晚上,实时总核销量应该等于离线核销量,差值超过千分之一就要报警。
对于差值,我一般用两个 SQL 去拉数。离线直接查 Hive 里的 DWS 汇总表,实时去查 ClickHouse。
-- 离线核销量 SELECT coupon_id, COUNT(*) AS off_cnt FROM dws_coupon_use_agg_daily WHERE dt = 'yyyy-MM-dd' GROUP BY coupon_id; -- 实时核销量 SELECT coupon_id, SUM(use_cnt) AS real_cnt FROM dws_coupon_use_agg WHERE day = 'yyyy-MM-dd' GROUP BY coupon_id;两张表一外关联,off_cnt和real_cnt不相等的记录就是问题点。导出差集后,针对差值大的券去 Kafka 里查它的明细 topic,看是不是 DWS 聚合窗口切错边界了。如果实时低于离线,大概率是 Flink 窗口终点设置为 ProcessingTime,而离线用的是 EventTime。具体表现是凌晨 0 点那批数据实时侧没有收进前一天的窗口,离线却算到了前一天,导致实时永远少一天。解决方式是改窗口 Flink 这边统一用事件时间戳来对齐,加上WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(10))。
这个对账技巧之所以强,是因为它逼着你去核对每条管道的数据血缘。我刚学会的时候,每天必调对账脚本,后来发现每一次差值都能指向一个具体的阶段,要么是 HBase 维表缺数据导致 join 丢了行,要么是 CH 写入超时导致漏掉一批。后来我把这个对账逻辑直接写进调度脚本,放在离线指标产出的前一个节点,每天自动跑一次,失败就从差值最大的 key 拆开下游链路深入排查。从那以后,我每次接新的数仓项目都会先问清楚实时和离线的最终结果表在哪,然后花一小时把对账脚本建好。排查问题的速度比别人快一倍,就是这笔投资换来的,希望帮到你。
本文还有配套的精品资源,点击获取