简介:围绕Spark离线数仓与Flink实时数仓双链路的项目源码与部署资料包,面向大数据开发学习者和求职者,完整覆盖实时数仓的ODS、DIM、DWD、DWS分层设计,以及离线场景的Spark批处理链路,直接解决项目实战和面试说理需求。资源共607个文件、约54.21MB,以Java源码、Class编译文件、SQL建表与数据脚本、Shell部署命令、XML/JSON配置为主,包含HBase、Kafka、ClickHouse等组件的实际接入代码和启动脚本,便于按层拆解实验。已有405人学习下载。实时数仓部分明确给出技术选型对比:ODS用Kafka逐条读取加工,DIM选HBase应对主键查询与长期存储,并解释了Redis内存有限、ES全字段索引、ClickHouse并发弱等原因;DWD继续用Kafka做分组累加,DWS落到ClickHouse进行聚合查询。通过源码结构、部署资料和分层笔记,可复用订单、用户、优惠券等业务域的实体类与Service实现,适合作为课程设计、求职项目或面试复盘的核心素材。
1. 拿到“Spark离线数仓Flink实时数仓项目源码+部署资料.rar”,先别急着解压
大多数人拿到这个压缩包的第一反应是解压、找 README、按部署文档敲命令,结果卡在版本不匹配、内存不够、任务起不来。这个标题背后是当前数据团队最主流的一套双轨架构:Spark 负责离线链路,从 Hive 分层加工出 T+1 报表;Flink 负责实时链路,用 CDC 接入业务库,经过实时 ETL 落到 Doris 或 Hive 供即时查询。它适合正在从单链路转向双轨的工程师,也适合准备面试时需要一个完整项目来拆解细节的从业者。下面要讲清楚两件事:代码怎么跑通、参数怎么调,以及哪些地方会悄无声息地翻车。和解压即跑不同,这个方向真正花时间的不是写代码,而是把部署资料里的环境约束和源码里的业务逻辑对上。
2. 双轨数仓的整体架构:离线数仓每一层的职责,与 Flink 实时链路的选型理由
2.1 离线数仓每一层的职责:Spark 在 ODS 到 ADS 之间到底做了什么
离线数仓每一层的职责在面试里是送分题,但在拿到源码后却容易被忽略。ODS 层保持和业务库一致,只做增量同步,不加工;DWD 层做清洗和维度退化,把 JSON 拍平、枚举翻译成中文;DWS 层按业务主题做汇总,比如订单主题、用户主题;ADS 层面向报表,输出最终指标。Spark 在这条链路里扮演的是批处理引擎,读 Hive 表、做宽表 Join、写回 Hive 分区表。市面上的网约车类大数据综合项目,基本都是这套分层再套一层业务外壳。
离线链路里最容易抄错的是 ODS 建表。很多新手把 ODS 直接建成内外表不分的普通表,后续 DWD 清洗时才发现分区策略根本没法覆盖重跑。按这类交付物的惯例,ODS 层应该用外部表按天分区,数据文件放 HDFS 的固定目录,这样即使业务表结构微调,也不会污染元数据。一个可供直接改用的 ODS 建表语句长这样:
CREATE TABLE IF NOT EXISTS ods.ods_order_detail_di ( order_id BIGINT COMMENT '订单ID', user_id BIGINT COMMENT '用户ID', sku_id BIGINT COMMENT '商品ID', num BIGINT COMMENT '购买数量', order_amount DECIMAL(10,2) COMMENT '订单金额', create_time TIMESTAMP COMMENT '下单时间' ) PARTITIONED BY (dt STRING COMMENT '分区字段,格式YYYY-MM-DD') STORED AS PARQUET;这个建表语句的关键点有三个:一是PARTITIONED BY (dt STRING),让离线任务天然按天调度,回刷历史数据时只需要删掉对应分区再重跑;二是STORED AS PARQUET,列式存储对后续聚合查询的扫描代价远低于 TEXT 格式,这也是 Spark 读 Hive 时最舒服的格式;三是 ODS 表不设 PRIMARY KEY,Hive 本身不支持主键约束,保持源系统的重复数据,去重逻辑延迟到 DWD 层做。这样设计后,离线数仓每一层的职责边界非常清楚,ETL 脚本之间不会互相踩踏。
2.2 实时数仓为什么选 Flink:事件时间与状态管理是硬门槛
实时链路选 Flink 而不是 Spark Streaming,不是因为流批一体的话术,而是三个具体理由:事件时间处理、原生状态管理、端到端精确一次。CDC 从 MySQL 采出来的数据,到达 Kafka 的时间顺序和业务库里的实际顺序不一定一致,Flink 的 Watermark 机制可以按create_time这样的业务时间戳来触发窗口计算,而不是依赖 Kafka 的到达时间。状态后端让维表关联、去重、窗口聚合这类算子把中间结果存在 RocksDB 或内存里,任务重启后还能从 checkpoint 恢复。Spark Streaming 做类似的事情需要把状态外置到 Redis 或 HBase,工程复杂度要高一个量级。
和入门时写的词频统计初体验不同,真实项目里的 Flink 实时计算几乎不会只处理一个无界流,而是要同时面对多个流和维表。源码里最常见的组织方式是:Flink SQL 定义 CDC 源表、Kafka 源表、Doris 维表,然后用一条 Insert Into 语句把整个实时 ETL 串起来。如果你发现业务数据不在 MySQL 而在自研消息队列里,就需要参考 Flink 实现自定义 Data Source 的方式,重写 RichParallelSourceFunction 来对接。选 Flink 的本质是选它的状态生态,而不是选一个能跑的流处理框架。
2.3 双轨在哪里汇合:Hive 表与实时结果表的数据服务口径
离线链路和实时链路不是平行的两条线,它们在 ADS 层会汇合。离线数据每天凌晨由 Spark 批量产出,写入 Hive ADS 表,供次日报表和领导驾驶舱使用;实时数据由 Flink 持续写入 Doris 明细或汇总表,供大屏、实时订单看板、异常监控使用。两条链路最终服务的是同一组业务指标,只是时效不同。这就产生了一个问题:同一指标在离线表和实时表里经常对不上,原因包括数据延迟、去重口径差异、Join 丢数据。部署资料里的绝大多数排查流程,最后都落在对账这一步。
所以在拆解源码时,建议先关注两条链路各自的结果表定义,再回头读中间层的清洗逻辑。一个很实用的习惯是,把离线 ADS 表的指标口径和实时 DWS 表的指标口径写成同一份文档,字段名、枚举值、时间口径都要一一对应。否则项目运行两周后,业务方拿着两张表来问为什么差了几万块钱,你会发现连排查的入口都找不到。压缩包里的部署资料如果带指标口径说明,这一份文档的价值往往比源码本身还高。
3. 拆开压缩包:源码目录结构与离线链路的最小复现路径
3.1 压缩包内的目录划分:按什么顺序读代码,才不会上来就迷路
这类交付物的目录通常按“部署资料 + 源码工程”两块组织。部署资料一般是 Markdown 或 PDF 文档,里面写了环境要求、集群规划、初始化脚本,以及从零跑通的步骤。源码工程则是 Maven 或 Gradle 工程,按模块拆成 spark-etl、flink-realtime、common 之类。一个容易踩的坑是:很多人一上来就打开 IDE 找主类,忽略了部署资料里的环境要求清单。Spark 和 Flink 对 JDK、Hadoop、Hive 的版本非常敏感,用 CDH 的集群跑纯 Apache 编译的代码,大概率会碰上依赖冲突。
我一般会按三条线去读这套源码。第一条线是 SQL 脚本目录,从 ODS 建表一路看到 ADS 建表,先把表结构和分区策略串起来;第二条线是 Spark 的 ETL 主类,确认每一层的输入输出表名,对照 SQL 脚本看加工逻辑是否一致;第三条线是 Flink 的实时任务,确认 CDC 源表的server-id、起始位点这类参数是不是硬编码。主线理清之后,再回头看部署资料里提到的资源规划,心里就有数了。按这个顺序读,两天之内能把这个压缩包里的业务逻辑吃透,而不只是跑通。
3.2 离线链路最小复现:用 Spark 读取业务 JSON 并写入 Hive DWD 表
解压源码之后,最有价值的操作是找出一条完整的离线 ETL 链路,从 ODS 到 DWD 先跑通。这里给一个贴近真实场景的模板:业务系统把订单数据以 JSON 格式落到 HDFS 目录,Spark 读取后做字段拍平、类型转换,再写入 Hive DWD 分区表。Spark 中读取 JSON 是离线分析最常见的入口之一,源码里对应的核心代码大致如下:
// 从HDFS读取业务系统的JSON文件 val rawDF = spark.read.json("hdfs://nameservice/data/raw/order_json/dt=2024-06-01") // 将嵌套JSON拍平,并把时间戳从字符串转成标准时间 val dwdDF = rawDF.selectExpr( "order_id", "user_id", "sku_id", "cast(amount as decimal(10,2)) as order_amount", "from_unixtime(ts, 'yyyy-MM-dd HH:mm:ss') as create_time", "get_json_object(ext_info, '$.province') as province" ) // 动态分区写入DWD层 dwdDF.write .mode("overwrite") .partitionBy("dt") .format("parquet") .saveAsTable("dwd.dwd_order_detail_di")这里的selectExpr是关键点。JSON 读取后天然是嵌套结构,直接用df.select("xxx")只能拿到顶层字段,嵌套对象里的 province 必须用get_json_object提取。from_unixtime(ts, 'yyyy-MM-dd HH:mm:ss')的作用是把秒级时间戳转成可读时间,很多项目在 ODS 层保留原始 ts,到 DWD 层才做格式化,这样历史分区回刷时不用重新解析原始文件。mode("overwrite")配合partitionBy("dt")时只覆盖分区目录而不动整张表,这是 Hive 表支持重跑的核心机制。
需要特别注意的是,saveAsTable默认会使用 Hive 的 SerDe 读取目录下的文件,如果之前用非 Parquet 格式写过同一个分区,会导致InputFormat不匹配的报错。因此 DWD 层的表必须在建表时就固定STORED AS PARQUET,一旦写入过其他格式,只能 drop 表重建。这个坑在回刷历史数据时几乎必现,压缩包里的部署资料如果没有单独提醒,你迟早会碰上。
3.3 spark-submit 提交参数:离线任务资源评估与三个必调项
离线链路跑通只是第一步,部署到集群上才是真正的考验。源码工程一般会带一个submit.sh脚本,里面是 spark-submit 的完整提交命令。以下是一份经过多次调参后比较稳妥的模板:
spark-submit \ --master yarn \ --deploy-mode cluster \ --name order_etl_daily \ --executor-memory 8g \ --executor-cores 4 \ --num-executors 20 \ --conf spark.sql.shuffle.partitions=200 \ --conf spark.sql.adaptive.enabled=true \ --conf spark.sql.adaptive.coalescePartitions.enabled=true \ --class com.example.etl.OrderEtl \ order-etl-1.0.jarspark.sql.shuffle.partitions是离线任务里最值得调的参数。默认值是 200,但如果你的 DWD 层数据量只有几 GB,200 个 reduce 任务会产生大量小文件;如果数据量有几百 GB,200 个分区又会造成单个任务处理时间过长。一般按“总数据量除以 128MB”估算分区数,再结合集群并行度取整。spark.sql.adaptive.enabled和coalescePartitions.enabled是 Spark 3 的 AQE 特性,开启后 Spark 会在运行时根据 shuffle 数据量自动合并小分区,配合spark.sql.adaptive.shuffle.targetPostShuffleInputSize一起用,能把小文件数量降一个量级。
Spark 内存问题也在这里暴露。executor-memory 8g是堆内内存,Spark 默认会把其中一部分留给 shuffle 和 storage,如果你发现任务频繁 Full GC,可以在 spark-defaults.conf 里调spark.memory.fraction=0.6,把堆内比例稍微调低。遇到 Executor Lost 之类的报错时,先查 YARN 日志里有没有“Container killed by YARN for exceeding memory limits”,这种几乎都是堆外内存超限,需要在spark-submit里额外加--conf spark.executor.memoryOverhead=2g。
4. Flink 实时链路:从 CDC 接入到维表关联与 Doris 落地
4.1 Flink CDC 安装部署与源表参数:binlog 开启是前置条件
实时链路的第一步是把 MySQL 的业务数据采进来。Flink CDC 的安装部署本身不复杂,但前置条件经常被忽略:MySQL 必须开启 binlog,并且 binlog 格式要设为 ROW。如果源库是阿里云 RDS 这类托管实例,binlog 的保留时长也要确认,否则消费位点过期后任务会从当前位置重新开始,造成数据断层。连接数据库的账号需要SELECT、RELOAD、SHOW DATABASES、REPLICATION SLAVE、REPLICATION CLIENT权限,缺一个权限 CDC 启动时就会报错。
Flink SQL 里定义一个 MySQL CDC 源表常见的写法是:
CREATE TABLE order_cdc ( order_id BIGINT, user_id BIGINT, sku_id BIGINT, num BIGINT, order_amount DECIMAL(10, 2), create_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = '10.10.1.20', 'port' = '3306', 'username' = 'cdc_user', 'password' = 'xxx', 'database-name' = 'trade_db', 'table-name' = 'order_detail', 'scan.startup.mode' = 'latest-offset', 'server-id' = '5400-5404' );server-id是 CDC 连接器最关键的参数。它用来标识 Flink 在 MySQL 复制链路里的从库身份,如果多个 CDC 任务共用同一个 server-id,MySQL 会判定为冲突连接,导致任务反复断开重连。给每个并发度分配一个独立的 server-id 区间是通行做法,比如并行度是 5,就用5400-5404的连续区间。scan.startup.mode有initial和latest-offset两个常用值。首次上线用initial,会先做全量快照再切到增量;日常重启用latest-offset,避免每次重启都重新扫全表。如果发现实时任务迟到了半小时以上,多半是initial同步全量时业务库压力过大,这个阶段要调大 MySQL 的max_allowed_packet和网络超时时间。
4.2 实时 ETL 里的维表关联:Lookup Join 的缓存与超时参数
实时链路里最影响吞吐的算子往往不是窗口聚合,而是维表关联。订单流里的sku_id要关联出商品名称、类目、品牌,用户 ID 要关联出省份城市。如果在 SQL 里用普通 Join,Flink 需要把维表全量加载到状态里,维表一大状态就膨胀。所以生产环境用的几乎都是 Lookup Join,也就是每条数据到达时实时查一次 Doris 或 HBase。在 Flink SQL 中定义一张 Doris 维表并关联的写法如下:
CREATE TABLE sku_dim ( sku_id BIGINT, sku_name STRING, category STRING, brand STRING, PRIMARY KEY (sku_id) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:mysql://10.10.1.30:9030/sku_db', 'table-name' = 'sku_dim', 'username' = 'dim_user', 'password' = 'xxx', 'lookup.cache' = 'PARTIAL', 'lookup.max-retries' = '3', 'lookup.cache.max-rows' = '10000', 'lookup.cache.ttl' = '1h' ); -- 实时订单流关联维表 CREATE VIEW order_with_sku AS SELECT o.order_id, o.order_amount, o.create_time, s.sku_name, s.category FROM order_cdc AS o LEFT JOIN sku_dim FOR SYSTEM_TIME AS OF o.proc_time AS s ON o.sku_id = s.sku_id;Lookup Join 的缓存参数直接决定维表查询对下游存储的压力。lookup.cache = 'PARTIAL'表示开启本地缓存,lookup.cache.max-rows控制缓存条数,lookup.cache.ttl控制缓存过期时间。这里最容易翻车的是把ttl设得太长——比如 24 小时。维表数据在业务库中每天都会更新,如果缓存 24 小时不刷新,实时报表里的商品名称就会隔天才能对上。反之,ttl设太短又会造成每条数据都穿透到 Doris,连接池被打满。一般取 30 分钟到 1 小时比较合适,同时要给维表所在的 Doris 或 MySQL 预留足够的连接数。
FOR SYSTEM_TIME AS OF o.proc_time是 Lookup Join 的固定语法,它的含义是“以订单数据的处理时间作为维表快照时间”。这个语法在 Flink 1.14 之后是强制要求,很多从旧版本迁移过来的 SQL 会在这里报语法错误。如果维表改成了 Flink CDC 维护的 HBase 表,语法不变,但连接器参数需要换成 HBase 的connector = 'hbase-2.2',同时指定zookeeper.quorum,缓存的调法完全相同。
4.3 实时结果落地:写入 Doris 的 Stream Load 参数与 checkpoint 配置
实时 ETL 的结果最终要服务查询。Doris 是这套架构里最常见的落地存储,既是实时结果表,也承担部分即席查询。Flink 写入 Doris 用的是 Stream Load 协议,在 SQL 里通过doris连接器声明一张结果表。一个供实时聚合结果的 Sink 表定义可以参考下面的写法:
CREATE TABLE doris_sku_agg_sink ( dt STRING, sku_id BIGINT, order_cnt BIGINT, amount_sum DECIMAL(14, 2) ) WITH ( 'connector' = 'doris', 'fenodes' = '10.10.1.30:8030', 'table.identifier' = 'dws.trade_sku_agg', 'username' = 'doris_user', 'password' = 'xxx', 'sink.label-prefix' = 'flink_doris_sku_agg', 'sink.properties.format' = 'json', 'sink.properties.read_json_by_line' = 'true', 'sink.enable.batch-mode' = 'true', 'sink.max-retries' = '3' );Doris Sink 的sink.label-prefix必须全局唯一。Stream Load 通过 label 保证导入幂等,如果两个任务共用同一个前缀,重启恢复时会把彼此的数据错认成同一次导入,产生脏数据。sink.properties.read_json_by_line要和实际写入的数据格式一致,默认 Stream Load 接收的是 CSV,改成 JSON 后每一行必须是一条完整 JSON,否则会报JSON Reader parse error。sink.enable.batch-mode是官方连接器推荐的批式写入模式,攒批再发能显著降低 Doris 的导入频率,对高峰期压力有明显缓解。
实时任务的可靠性最终由 checkpoint 兜底。Flink 默认的 checkpoint 间隔偏保守,生产环境里我会把execution.checkpointing.interval从默认值调到 60 秒,execution.checkpointing.min-pause设成 30 秒,同时开启execution.checkpointing.externalized-checkpoint-cancel为 true。这样做的目的是让任务在停止或故障后,能够从最近一次成功的 checkpoint 恢复,配合 Doris 的 label 幂等机制,能做到“任务重启不重账”。如果你发现恢复后数据翻倍,问题不在 checkpoint,而是 Sink 端没有幂等保证,这时优先检查sink.label-prefix是否在任务重启前后保持一致。
5. 避坑记录:部署与运行阶段最常见的九类翻车现场
5.1 环境与版本类:部署资料在这套集群能跑,换一套就起不来
现象:严格按照部署资料执行,Hadoop、Hive 版本一致,Spark 任务却在提交后秒失败,报错信息大量出现ClassNotFoundException或NoSuchMethodError。原因不是集群坏了,而是部署资料里的版本和实际环境存在隐藏差异,最常见的有三种:一是 Spark 用 Scala 2.12 编译,集群上的 Spark 是 Scala 2.11 版本;二是 Hive 的 metastore 版本不同导致HiveConf初始化失败;三是 Hadoop 的mapreduce相关依赖冲突。
解决:拿到部署资料后先核对三个版本号——spark.version、hadoop.version、scala.version。源码工程里的pom.xml或build.sbt会有明确声明,集群上通过spark-submit --version和hadoop version直接查。版本不一致时优先改源码重编译,而不是试图替换集群里的原生 jar。用mvn clean package -DskipTests -Pdist重打一次包通常能解决九成环境问题。如果重编后还报依赖冲突,就在 spark-submit 里加--conf spark.driver.userClassPathFirst=true和--conf spark.executor.userClassPathFirst=true,把应用自带的依赖优先加载。
5.2 Flink 提交任务后 TaskManager 一直重启:不只是内存不够
现象:Flink 任务提交后 JobManager 正常,但 TaskManager 反复加入又退出,YARN 日志里能看到Container killed by YARN for exceeding memory limits,甚至直接Exit code: 143。新手第一反应是调大容器内存,结果开得越大,被 kill 得越快。原因通常是 Flink 进程的堆外内存超限,taskmanager.memory.process.size设得很大,但jvm-overhead没有随之调整,RocksDB 状态后端的堆外占用把容器顶爆了。
解决:先明确 Flink 的完整内存模型,再调参数。taskmanager.memory.process.size是总内存,它包含堆内、堆外、JVM Overhead、网络缓冲四部分。一个常用的起点是:总内存给 8g,taskmanager.memory.managed.size=512m(RocksDB 跑时再按需加),taskmanager.memory.jvm-overhead.fraction=0.2。如果用了 RocksDB,还要注意state.backend.rocksdb.memory.managed=true,让 RocksDB 使用受管内存而不是无限制吃堆外。调完参数后让任务空跑半天盯着监控,确认内存曲线平稳后再接真实流量。
5.3 实时链路数据漂移:Flink JDBC 连接器异常与 Hive 表写不进数据
现象一:Flink 任务跑了一周后突然报Communications link failure,重试几次后恢复,但期间实时报表出现分钟级空窗。检查后发现是 JDBC 连接器连接池里的连接被 MySQL 服务端主动断开,而 Flink 没有自动重建。解决办法是在连接器参数中长期维护一批连接池参数:jdbc.connection.pool.size=10、jdbc.connection.max-retry-timeout=60,同时把 MySQL 端的wait_timeout调到不低于连接池的maxLifetime。
现象二:Flink 任务正常结束,但 Hive 结果表里始终查不到数据。原因多半是对flink sink hive的机制理解有偏差:Hive Sink 不是实时写入文件,而是通过streaming模式 + 分区提交来可见数据。没有设置sink.partition-commit.policy.kind=success时,即使数据写入了 HDFS,Hive 表也读不到。解决方法是给 Sink 表加上分区提交参数,并设置sink.partition-commit.trigger=process-time。注意 Hive 表只能用动态分区写入,分区字段必须在 SELECT 里显式携带,否则默默全量覆盖。
5.4 离线统计值异常:Spark 读 Hive 分区表的数据量忽大忽小
现象:离线 ADS 表每天跑出来的订单总量偶尔比前一天少几万条,反复重跑数字却稳定不变,说明不是计算抖动,而是上游分区数据出了问题。最常见的原因是 ODS 层的增量文件在当天推送时发生过覆盖,Spark 读到的是半新半旧的快照。另一个原因是 DWD 层在动态分区写入时产生了大量小于 1MB 的小文件,后续读取时spark.sql.files.maxPartitionBytes默认 128MB 对小文件不敏感,但 NameNode 压力会变大,任务耗时明显变长。
解决:把 ODS 到 DWD 的读取操作加上版本校验,最简单的做法是在 ODS 表里增加一个etl_time字段,DWD 只取当天最大etl_time的数据。对已经产生的小文件,用spark.sql.adaptive.coalescePartitions.enabled=true配合spark.sql.adaptive.advisoryPartitionSizeInBytes=64MB重新整理。如果业务上允许,定期对 DWD 分区执行一次INSERT OVERWRITE,用distribute by dt打散文件分布,比任何参数都管用。
5.5 对账对不上:离线 Hive 与实时 Doris 的同一指标差了几千单
现象:离线 ADS 表和实时 Doris 大屏上的“今日订单数”差了 3000 多单,两边单独看都“对”,放一起就对不上。排查后发现问题出在两个地方:第一,离线表统计的是create_time在当天的订单,实时表统计的是 Flink 处理时间在当天的订单,延迟数据跨天导致差异;第二,实时链路把同一订单的更新操作(比如状态从“已支付”改成“已发货”)算成了两单,而离线表在 DWD 层做了去重。解决方式很简单:把两张表的口径统一成“基于create_time的订单唯一键去重”,实时 DWS 层增加COUNT(DISTINCT order_id)的窗口计算即可。对账时优先用唯一键逐条比对,而不是先比聚合值。
6. 用一张对账单验证双轨数据一致性:离线 ADS 与实时 Doris 的差值核查技巧
双轨架构上线后最该做的事不是加更多报表,而是把离线 Hive 与实时 Doris 的对账做成每天自动执行的任务。我常用的做法是写一条 SQL,把同一指标的两份结果放在一张表里做减法,差值超过阈值就告警。以“当日订单金额”为例,核心 SQL 长这样:
SELECT sku_id, MAX(offline_amount) AS offline_amount, MAX(realtime_amount) AS realtime_amount, MAX(offline_amount) - MAX(realtime_amount) AS diff FROM ( SELECT sku_id, amount_sum AS offline_amount, 0 AS realtime_amount FROM hive_catalog.ads.trade_sku_ads WHERE dt = '2024-06-01' UNION ALL SELECT sku_id, 0 AS offline_amount, amount_sum AS realtime_amount FROM doris_catalog.dws.trade_sku_agg WHERE dt = '2024-06-01' ) t GROUP BY sku_id HAVING ABS(MAX(offline_amount) - MAX(realtime_amount)) > 100;这条 SQL 的原理是把两张表的数据按sku_id横向拼接成一行,再通过MAX把两个值取出来做差。HAVING ABS(...) > 100直接过滤出差异超过阈值的商品,每次只需要关注这几十条记录即可。跑完 SQL 后,把结果接入一个简单的告警脚本,每天凌晨五点执行一次,只有存在 diff 时才推送告警。对账脚本的价值不在于多复杂,而在于能坚持跑。很多项目上线初期对账一切正常,三个月后业务表变更频繁,问题才会集中爆发。
我个人的习惯是,每次交付双轨数仓项目,都会先在部署资料里补齐一份对账 SQL 模板。原因是离线表和实时表的数据口径太容易漂移了,靠人肉盯报表根本不现实。把对账做成例行任务后,你会发现自己对两套链路的健康状况比业务方还清楚。这个经验可以平移到你自己的项目里,但也别过度依赖——对账只负责发现问题,真正修数据要靠状态恢复和分区重跑,两种手段配合才能把双轨数仓守住。希望帮到你。
本文还有配套的精品资源,点击获取