惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

博客园 - 司徒正美
Jina AI
Jina AI
Microsoft Azure Blog
Microsoft Azure Blog
博客园 - 三生石上(FineUI控件)
宝玉的分享
宝玉的分享
MyScale Blog
MyScale Blog
I
InfoQ
爱范儿
爱范儿
Microsoft Security Blog
Microsoft Security Blog
酷 壳 – CoolShell
酷 壳 – CoolShell
Stack Overflow Blog
Stack Overflow Blog
T
Tailwind CSS Blog
D
DataBreaches.Net
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
T
The Blog of Author Tim Ferriss
B
Blog
阮一峰的网络日志
阮一峰的网络日志
Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC
月光博客
月光博客
雷峰网
雷峰网
Recent Announcements
Recent Announcements
量子位
B
Blog RSS Feed

博客园 - AlfredZhao

Oracle Property Graph 结合 In-Memory 的性能验证 sed 一条命令批量替换 SQL 文件路径 执行新项目 python 脚本前,先用 conda 建一个独立环境 Git 提交代码:从报错到 SSH 免密推送 GitHub 克隆他人私有仓库:从授权到下载 理解Oracle Property Graph:以账户转账示例完成图特性最小测试 APEX 无法分配 SH 用户?一个存储过程轻松解决 SH 中文化样例数据使用手册 Oracle 排除非业务表:一份能直接抄的“全量过滤”SQL 进程都杀了,为什么 `netstat` 还能看到端口? MAC 空间告急?Codex 的 156GB 缓存垃圾,这样清! 人工清理问题数据:先查准,再删除 TK(Trusted Knowledge)为何而生? 使用快捷键快速切换 Mac 外接显示器模式 客户环境 Nginx 配置:流式报表与超时排查要点 Security Central:数据库安全的统一控制与运营平台 Palantir 眼中的一次“订单可能延期”,如何成为实时决策的起点? UUID v4 与 v7:同样是唯一 ID,为什么数据库表现可能完全不同? 用 crontab 给 LLM 使用量装上“监控眼” DeepSeek 与 GPT API 价格调整:该关注什么? 查看 Oracle 数据库中的定时任务执行情况 Codex 专用用户登录后自动进入默认项目目录 Mac 小技巧:用方向键优雅处理超长网址 Oracle GDD 与 Raft:别把共识机制和数据主权混为一谈 在 Oracle APEX 中用 BGE_BASE 生成向量:从模型导入到历史数据更新 行业人+AI:真正有价值的四个关键要素 知识库文件解析失败:一次由 Domain Index 引发的定位记录 Unicode 码位数、UTF-8 字节数、中英文差异和 Oracle 长度别再混淆了 Skill 的使用:从路径到能力边界 切换 Embedding 模型时,千万别忽略历史向量维度
从一条订单消息到 Bronze、Silver、Gold:我的第一次 Kafka +...
AlfredZhao · 2026-09-01 · via 博客园 - AlfredZhao

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

配套代码和完整运行步骤,后续如有读者感兴趣再单独写一篇文章分享。

01 | Kafka 传递事件,不替代数据库

订单表通常描述“当前状态”:例如订单 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 服务事件分发与处理。

02 | Key、Partition 与 Offset 的边界

实验创建了 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_scnkafka_offset 也不是同一件事:

source_scn   → 模拟 Oracle 源端的变更位置
kafka_offset → 消息在 Kafka Partition 中的位置

CDC 链路中,两者都值得保留。

03 | Bronze、Silver、Gold 分别解决什么

Consumer 将 Kafka 消息写入:

lake/bronze/order_events.jsonl

每条记录附带 kafka_topickafka_partitionkafka_offsetingested_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 等业务场景。

04 | 血缘让链路可追溯

实验生成 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。发生金额异常或字段变更时,也能沿血缘进行影响分析和排查。

05 | 最小实验的价值

当前实现使用 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一起成长~