














2026-09-01 06:53 AlfredZhao 阅读(0) 评论() 收藏 举报
作为 Oracle DBA,笔者熟悉表、事务、Redo、SCN 和数据仓库;但第一次面对 Kafka、Topic、Partition、Offset 与 Bronze、Silver、Gold 时,仍觉得概念很多、关系不清。于是笔者用 6 条订单事件跑通了一条最小链路。
graph TD Oracle([模拟 Oracle CDC 订单事件]) KafkaTopic["Kafka Topic: shop.orders (Key: order_id | 3 Partitions)"] subgraph Bronze [Bronze 原始证据层] F_Bronze[lake/bronze/order_events.jsonl] Note_Bronze[保留全部 6 条历史状态变化] end subgraph Silver [Silver 可信明细层] F_Silver[lake/silver/orders.parquet] Note_Silver[DuckDB 按 order_id 分组去重] end subgraph Gold [Gold 业务指标层] F_Gold[lake/gold/daily_order_kpi.parquet] Note_Gold[聚合汇总: 日订单 KPI] end Lineage[(数据血缘: lineage.json)] Oracle --> KafkaTopic KafkaTopic -->|consume_to_bronze| F_Bronze F_Bronze -->|deduplicate_and_validate_orders| F_Silver F_Silver -->|aggregate_daily_order_kpi| F_Gold F_Bronze -.-> Lineage F_Silver -.-> Lineage F_Gold -.-> Lineage
配套代码和完整运行步骤,后续如有读者感兴趣再单独写一篇文章分享。
订单表通常描述“当前状态”:例如订单 1001 当前为 PAID。
事件流更关注过程:
graph LR subgraph DB [Oracle 传统关系型数据库] direction TB Table[(订单表: orders)] State["当前状态:订单 1001 = PAID"] Table --- State end subgraph Kafka [Kafka 分布式事件流] direction TB Log[["Kafka Topic: shop.orders (追加日志)"]] E1["1. 订单 1001 已创建"] --> E2["2. 订单 1001 已付款"] --> E3["3. 订单 1001 已发货"] Log --- E1 end style DB fill:#f4f5f7,stroke:#6b778c; style Kafka fill:#fff0f0,stroke:#de350b;
Kafka 是面向业务事件的分发通道,可供多个消费者订阅和回放。笔者暂时将它理解为分布式追加日志:它与 Oracle Redo 都有位置和顺序概念,但 Redo 服务恢复与一致性,Kafka 服务事件分发与处理。
实验创建了 3 个 Partition 的 shop.orders Topic,并以 order_id 作为消息 Key:
producer.produce(
"shop.orders",
key=str(item["order_id"]),
value=json.dumps(item),
)
同一订单的事件会进入同一 Partition,因此其状态顺序稳定;不同订单则可能因哈希落入同一 Partition。
Kafka 保证的是同一 Topic-Partition 内的顺序,不保证整个 Topic 的全局顺序。因此,不能根据不同 Partition 的消费输出先后判断业务发生顺序。
source_scn 与 kafka_offset 也不是同一件事:
source_scn → 模拟 Oracle 源端的变更位置
kafka_offset → 消息在 Kafka Partition 中的位置
CDC 链路中,两者都值得保留。
Consumer 将 Kafka 消息写入:
lake/bronze/order_events.jsonl
每条记录附带 kafka_topic、kafka_partition、kafka_offset 和 ingested_at。Bronze 保存全部 6 条状态变化,而非每个订单的最新状态。
Bronze 是原始证据层:规则出错、口径变化或需要重算时,仍能回到原始事件。
Silver 使用 DuckDB 按 order_id 分组,并依 source_scn desc, kafka_offset desc 选择最新记录,最终得到 3 个订单的可信明细:
Bronze 回答:历史上发生过什么?
Silver 回答:经过规则处理后,目前可信的明细是什么?
Gold 则将可信明细转为业务指标。实验结果中,3 个订单有 2 个已付款、1 个已取消,付款金额为 135.40:
99.90 + 35.50 = 135.40
Gold 不再关注 Partition 或 Offset,而是服务报表、运营、API 或 AI 等业务场景。
实验生成 lake/lineage.json,记录加工路径:
kafka://shop.orders
↓ consume_to_bronze
lake/bronze/order_events.jsonl
↓ deduplicate_and_validate_orders
lake/silver/orders.parquet
↓ aggregate_daily_order_kpi
lake/gold/daily_order_kpi.parquet
这使笔者能够回答:某个 Gold 指标来自哪里、经过哪些转换,以及源头是哪一个 Kafka Topic。发生金额异常或字段变更时,也能沿血缘进行影响分析和排查。
当前实现使用 Kafka、JSONL、Parquet、DuckDB 和静态血缘文件,主要演示 Lakehouse 的分层思路,并非生产级湖仓平台。
后续可逐步将模拟 Producer 替换为 Oracle Free + Debezium LogMiner CDC,将本地目录替换为 OCI Object Storage 或 MinIO,并引入 Apache Iceberg、Flink SQL 与 OpenLineage。
这次实验只有 6 条消息,却让笔者真正看见了事件如何进入 Partition、Offset 如何增长、Bronze 如何保留证据、Silver 如何按 SCN 去重,以及 Gold 如何形成业务指标。对 Oracle DBA 而言,事务、日志、SCN、分区和分析函数等经验,正是理解 Kafka 与 Lakehouse 的桥梁。
关注我,和AI一起成长~
此内容由惯性聚合(RSS阅读器)自动聚合整理,仅供阅读参考。 原文来自 — 版权归原作者所有。