Data warehouse / Doris / Flink / Paimon

电商离线 + 实时
双版本湖仓数仓

同一组电商交易指标,分别通过 Doris 离线物理分层和 Flink SQL 实时逻辑计算链产出;实时结果并行物化到 Paimon 分层表与 Doris ADS,最终在 Doris 统一查询并逐字段对账。

为什么要做双链路数仓

离线链路强调稳定重算和历史口径,实时链路强调持续更新和低延迟看板。真正的难点不是让两个系统各跑一遍,而是让同一业务指标在两套技术栈里有一致口径、可解释结果和可复现验收。

指标口径容易分裂

离线报表和实时看板如果各自实现,GMV、订单数、退款订单数很容易因为过滤条件或日期范围不同而不一致。

技术职责需要拆清楚

Doris 承担 OLAP 查询服务和统一 ADS 出口;Paimon 保存同一实时逻辑 DAG 并行物化出的 ODS/DWD/DWS/ADS 结果。

聚焦可复现的核心闭环

本次实现聚焦 MySQL、Doris、Kafka、Flink 与 Paimon,重点完成离线与实时双链路的本地复现和逐字段对账,避免额外组件稀释核心验证目标。

整体架构:同指标两套实现

离线链路从 MySQL 业务库进入 Doris 并按物理表逐层加工;实时链路从 Kafka 事件进入 Flink SQL,由临时视图组成 DWD → DWS → ADS 逻辑计算链,再通过同一 Statement Set 并行写入 Paimon 各层与 Doris ADS。最后在 Doris 对账 offline ADS 与 realtime ADS。

01

Doris 离线数仓

MySQL 同步到 Doris ODS/DIM,再通过 Doris SQL 加工到 DWD、DWS、ADS。

02

Flink + Paimon 实时湖仓

Kafka 事件进入 Flink SQL 临时视图 DAG,Statement Set 将 ODS/DWD/DWS/ADS 与 Doris ADS 作为多个 sink 并行物化。

03

Doris 统一查询

对齐 `dt + recent_days` 粒度,逐字段比较 GMV、去重订单/用户、退款订单/用户及客单价。

电商离线物理分层与实时逻辑视图链并行物化架构图
架构图突出三条边界:MySQL 到 Doris 的离线物理分层、Kafka 进入 Flink SQL 内部逻辑视图链后并行写入 Paimon/Doris、Doris 中 offline/realtime ADS 对账。

Doris 离线数仓链路

离线链路用 MySQL 模拟电商业务库,脚本同步到 Doris ODS 和 DIM,再通过 Doris SQL 完成 DWD 明细标准化、DWS 固定日期范围汇总和 ADS 指标产出。

MySQL ecommerce_oltp -> Doris ODS / DIM -> Doris DWD -> Doris DWS -> ads.ads_trade_stats_offline
Doris ODS DIM DWD DWS ADS 分层表行数查询结果
Doris 查询一次性列出 ODS、DIM、DWD、DWS、ADS 的表行数,证明离线链路不是只落最终 ADS,而是按数仓分层产出可检查的中间结果。

Flink + Paimon 实时湖仓链路

实时链路使用 Kafka 确定性事件作为输入。DWD、DWS、ADS 由 Flink SQL 临时视图串成逻辑计算链;BEGIN STATEMENT SET 中的多个 INSERT 再把 ODS、DWD、DWS、ADS 并行物化到 Paimon,并通过 Doris Flink Connector 同步物化最终 ADS。Paimon 下游表在这个作业中不是逐层读取上游 Paimon 物理表。

Flink Dashboard 实际运行截图,作业状态为 RUNNING,9 个 task 的 DAG 和算子状态均可见
Flink Dashboard 实际运行视图:顶部 Job State 为 RUNNING,作业包含 9 个 task;DAG 与下方算子表同时展示 Kafka Source、DWD/DWS 聚合以及 Paimon/Doris 多个 writer。最终 fresh clone 的 checkpoint 与落盘验收见“本地验收入口”。

指标合同:同一口径才有对账意义

两条链路都按 `dt + recent_days` 输出,核心指标包含 GMV、去重订单数、去重下单用户数、退款订单数、退款用户数和客单价。业务日期固定为 2026-07-01,`recent_days` 表示固定日期范围,并非 Flink Window TVF。

recent_daysGMV订单数下单用户数退款数退款用户数客单价
1505.505422101.10
71078.509633119.83
301578.5010744157.85

Doris 统一查询与逐字段对账

最终不是只看两个表都有数据,而是在 Doris 中把 offline ADS 和 realtime ADS 按 `dt + recent_days` 对齐,用同一次查询逐字段比较 GMV、去重订单/用户、退款订单/用户和客单价。三个固定统计范围全部 PASS,说明两条链路在当前样例数据和指标合同下结果一致。

Doris Playground 对账 SQL 结果截图,1、7、30 日三个固定统计范围均为 PASS
Doris Playground 直接对齐 offline ADS 与 realtime ADS,1 / 7 / 30 日三个固定统计范围均返回 PASS;完整指标口径与最终 3 / 3 验收结果分别记录在上方指标合同和“本地验收入口”。

关键工程取舍:少堆组件,多证明闭环

Doris 作为统一查询出口

离线 ADS 和实时 ADS 都写入 Doris,BI 或演示只需要查一个 OLAP 层,也能直接做 offline/realtime 对账。

逻辑计算与物化解耦

DWD → DWS → ADS 的依赖由临时视图表达;Paimon ODS/DWD/DWS/ADS 和 Doris ADS 是同一 Statement Set 的并行物化目标,用于检查、留存与统一查询。

Doris Connector 替代 JDBC sink

通用 JDBC sink 会遇到 Doris MySQL 协议与 `ON DUPLICATE KEY` 语法兼容问题,因此改用 Doris Flink Connector 走 Stream Load。

确定性样例保证可复现

Kafka topic 每次重置后再写入样例事件,避免重复生产导致实时 GMV 翻倍,保证离线和实时可以稳定对账。

本地验收入口

第一版以 Docker Desktop + WSL2 为本地运行环境,启动 MySQL、Doris、Kafka、Flink 后按脚本顺序执行即可复现。

docker compose ps 截图,展示 MySQL、Doris、Kafka、Flink 服务运行
本地服务启动后,Doris、MySQL、Kafka、Flink JobManager 和 TaskManager 均处于运行状态。
Fresh clone PASS · 5 services · Flink RUNNING · 8 RUNNING + 1 FINISHED task · checkpoint 1 → 2 with no new failure · Paimon 5 / 5 current-job files · 3 / 3 reconciliation PASS
powershell -NoProfile -ExecutionPolicy Bypass -File .\scripts\validate-repo.ps1
powershell -NoProfile -ExecutionPolicy Bypass -File .\scripts\run-demo.ps1 -Reset

当前边界与增强路径

这版的价值是跑通双链路指标合同和本地复现闭环,不把尚未完成的增强项写成已交付能力。

当前边界

  • 实时输入为 Kafka 确定性样例事件,MySQL CDC 尚未接入。
  • DolphinScheduler 尚未编排,当前用 PowerShell 脚本串联。
  • Paimon 使用本地 filesystem catalog,未接 Hive Metastore 或对象存储。
  • 实时分层是临时视图逻辑链加并行 sink,不是 Paimon 物理表之间逐层读取。
  • 1 / 7 / 30 日指标采用固定日期条件持续聚合,尚未实现 Window TVF、乱序与迟到数据处理。
  • 样例数据规模小,重点是口径对齐和链路闭环,不是压测性能。

下一步增强

  • 接入 DolphinScheduler,把离线初始化、同步和 ADS 聚合编排成 DAG。
  • 增加 MySQL CDC 或 Debezium/Flink CDC 输入。
  • 增加 Window TVF、乱序和迟到数据用例,补齐事件时间语义。
  • 把端到端容器验收接入可重复运行的专用 CI 环境。

这个项目想证明什么

我能把同一个电商交易指标拆成 Doris 离线物理分层与 Flink SQL 实时逻辑计算链两套实现,使用 Statement Set 并行物化 Paimon 分层结果和 Doris ADS,再通过 Doris 逐字段对账;同时能说明本地 MVP 的实现边界。