raiza-king/rk-medallion-arc-proj
GitHub: raiza-king/rk-medallion-arc-proj
该项目将新南威尔士州的实时燃油价格和年度交通流量数据,通过 Databricks 上的 medallion 幂等流水线整合为一个面向燃油预算和货运规划的决策支持仪表板。
Stars: 0 | Forks: 0
## 1. 业务背景
燃油是物流中最大的*可变*成本,且货物所经公路的路况同时决定了配送时间和路线成本。将**实时**燃油价格数据流与**战略性**历史交通数据流结合,能够讲述任何单一数据源都无法单独展现的完整故事。
**目标用户:**管理供应链和预算的中小型企业,以及需要为 NSW 偏远地区的燃油做预算的个人和家庭。
**它回答的两个问题:**
- **预算规划** — 燃油价格走势如何?下个月我该预留多少资金?
- **路线规划** — 哪些走廊承载的货运量最大?沿线的燃油价格情况如何?
## 2. 数据源
| 数据源 | 内容 | 频率 | 访问方式 |
| --- | --- | --- | --- |
| **NSW FuelCheck / Fuel API** | 每个 NSW 注册加油站的实时燃油价格 | 实时 / 运营 | 2 步 OAuth + 必需的 headers |
| **TfNSW Traffic Volume Viewer** | 按车型、方向和时段统计的年度车辆计数(自 2006 年起约 1,900 个站点) | 年度 / 战略 | 批量 CSV,手动上传 |
## 3. 架构
为 **Serverless + Unity Catalog** 工作区量身定制的经典 **medallion** 架构——无需看管任何 cluster,无需管理任何存储账号或访问密钥。原始文件落地到 UC 管理的 **volume** 中;silver 和 gold 层是受管的 **Delta tables**。
```
flowchart LR
FUEL["NSW FuelCheck API
(live prices)"] --> BRONZE TRAFFIC["TfNSW Traffic CSV
(yearly counts)"] --> BRONZE BRONZE["Bronze
raw files (UC volume)"] --> SILVER SILVER["Silver
clean, typed, de-duped
Delta tables"] --> GOLD GOLD["Gold
business aggregates"] --> DASH["SQL Dashboard
budgeting + logistics"] ``` | 层级 | 目的 | 存储形式 | | --- | --- | --- | | 🥉 Bronze | 原始、未修改的摄入数据,完全保持拉取时的原貌 | UC volume 中的文件 | | 🥈 Silver | 已清洗、已转换格式、去重且标准化的数据 | 受管的 Delta tables | | 🥇 Gold | 适用于仪表板的业务级聚合数据 | 受管的 Delta tables | ## 4. Notebook 结构 该流水线存在于单个 notebook 中,按顺序划分为多个部分: | 部分 | 作用 | | --- | --- | | `00_manualcleanup` | **仅供开发 / 恢复。** 清空落地收件箱,全量重建删除,一次性丢弃废弃表,以及安全的 gold 裁剪脚本。从不进行调度。 | | `01_bronze` | 创建 catalog / schemas / volume,然后拉取原始燃油 (API) + 交通 (CSV) 数据。仅采用追加模式。 | | `02_silver` | 通过 **MERGE** 将 bronze 转换为 silver(幂等 upsert、去重、格式化)。 | | `03_gold` | 通过确定性的 `overwrite` 从 silver 层构建业务聚合数据。 | ## 5. 幂等性设计(核心原则) 该流水线的设计确保了**重新运行任何层级都是安全的,并且会产生相同的结果**——没有重复,也不会发生意外数据丢失。 - **Bronze 层仅支持追加。** 每次 API 拉取都会写入一个带有新时间戳的文件 (`fuel_.json`),从而保留不可变的原始审计跟踪。常规删除操作已被移除。
- **Silver 层仅通过 MERGE 摄入新行或已更改的行。** 燃油价格仅在 `(station_code, fuel_type, updated_at)` 上执行插入;加油站基于 `station_code` 执行 upsert;交通数据则按其自然粒度执行 merge。重新运行绝不会导致观测数据重复。
- **Gold 层完全派生,并使用 `overwrite` / `CREATE OR REPLACE` 进行重建。** 由于它是基于 silver 确定性生成的,因此全量重建是幂等的,并且只会触及目标 gold 表。
- **完全重置已被隔离。** “删除所有内容并重新运行”的操作位于 `00_manualcleanup` 下,专用于 schema 重设计、状态损坏或受控的数据回填。
## 6. Gold 表与血缘
仪表板直接读取三个 gold 表。繁重的中间计算过程是在 `03_gold` 运行期间**在内存中**构建的(不持久化),因此重新运行是完全独立的,也不需要维护孤立的表。
| 表名 | 是否由仪表板读取? | 备注 |
| --- | --- | --- |
| `gold.fuel_price_trends` | ✅ 直接读取 | 每种燃油类型每小时的平均/最低/最高价格(预算时间序列) |
| `gold.servo_corridor` | ✅ 直接读取 | 每个加油站与其最近的货运走廊配对(Haversine 空间 join) |
| `gold.map_layer` | ✅ 直接读取 | 加油站和走廊作为独立点位堆叠,用于在同一个合并地图中展示 |
| `gold.heavy_daytype_by_corridor` | ✅ 直接读取 | 仪表板数据集(受保护) |
| `servo_map_df` / `corridor_load_df` / `heavy_by_period_df` | ⚙️ 内存中 | 在 runtime 期间计算的中间变量;**不**作为表持久化 |
## 7. 如何运行
**前置条件**
- 启用了 Unity Catalog 的 Databricks Serverless 工作区
- NSW Fuel API key + secret 存储在名为 `medallion` 的 secret scope 中(键名:`nsw_fuel_key`, `nsw_fuel_secret`)
- Traffic Volume Viewer CSV 已上传至 `medallion.bronze.landing/traffic`
**执行顺序**
1. 运行 `01_bronze` 一次,以创建 catalog / schemas / volume,随后拉取原始数据。
2. 运行 `02_silver` 以构建已清洗的 Delta tables(可安全地每日重跑)。
3. 运行 `03_gold` 以构建仪表板聚合数据。
4. 将 serverless SQL warehouse 或仪表板指向这些 gold 表。
## 8. 维护与安全防护 (`00_manualcleanup`)
- **清理落地收件箱** — 仅用于开发/测试;否则的话,bronze 层的原始历史记录会被永久保留。
- **彻底重建 silver/gold 层** — 仅用于恢复或 schema 重设计。
- **一次性丢弃废弃表** — 移除旧的已物化中间表 (`servo_fuel_map`, `freight_corridor_load`, `heavy_by_period`),它们现在都在内存中处理。
- **Gold 裁剪(允许列表 + dry-run)** — 列出一个 `PROTECTED` 仪表板表保留集合,并丢弃任何其他 gold 表。默认设置为 `DRY_RUN = True`,以便您在删除前进行预览。
## 9. 仪表板
一个针对两类用户角色设计的 Databricks SQL 仪表板,通过 AI/BI 嵌入式 SDK,借助部署在 Azure App Service(免费层)上的小型 Flask + React 应用嵌入到 Databricks 之外。
- **💰 预算视图** — 随时间变化的燃油成本指数、当前价格计数器、最便宜品牌的柱状图。
- **🚚 物流视图** — 带有价格和最近走廊工具提示的加油站地图,以及各时段的货运量。
## 10. 技术栈
Databricks Serverless · Unity Catalog · PySpark · Delta Lake · Databricks SQL · Flask + React (嵌入式) · Azure App Service。
## 11. 治理与后续计划
如果该项目从个人构建成果转变为共享产品,我们将采取以下措施:采用最小权限的 UC 访问控制、dev/test/prod 环境分离、明确数据集所有者及数据新鲜度 SLAs、实施数据质量期望、完善 UC 血缘关系与表注释、对 gold 层逻辑实行 PR 审查,并针对 serverless 的使用情况设置成本防护机制。
**延伸目标:** 燃油支出场景计算器(升/月 × 柴油价格,±10¢ 敏感度分析)、向 Delta Live Tables / Lakeflow 的迁移,以及扩展至墨尔本 (VIC) 地区。
(live prices)"] --> BRONZE TRAFFIC["TfNSW Traffic CSV
(yearly counts)"] --> BRONZE BRONZE["Bronze
raw files (UC volume)"] --> SILVER SILVER["Silver
clean, typed, de-duped
Delta tables"] --> GOLD GOLD["Gold
business aggregates"] --> DASH["SQL Dashboard
budgeting + logistics"] ``` | 层级 | 目的 | 存储形式 | | --- | --- | --- | | 🥉 Bronze | 原始、未修改的摄入数据,完全保持拉取时的原貌 | UC volume 中的文件 | | 🥈 Silver | 已清洗、已转换格式、去重且标准化的数据 | 受管的 Delta tables | | 🥇 Gold | 适用于仪表板的业务级聚合数据 | 受管的 Delta tables | ## 4. Notebook 结构 该流水线存在于单个 notebook 中,按顺序划分为多个部分: | 部分 | 作用 | | --- | --- | | `00_manualcleanup` | **仅供开发 / 恢复。** 清空落地收件箱,全量重建删除,一次性丢弃废弃表,以及安全的 gold 裁剪脚本。从不进行调度。 | | `01_bronze` | 创建 catalog / schemas / volume,然后拉取原始燃油 (API) + 交通 (CSV) 数据。仅采用追加模式。 | | `02_silver` | 通过 **MERGE** 将 bronze 转换为 silver(幂等 upsert、去重、格式化)。 | | `03_gold` | 通过确定性的 `overwrite` 从 silver 层构建业务聚合数据。 | ## 5. 幂等性设计(核心原则) 该流水线的设计确保了**重新运行任何层级都是安全的,并且会产生相同的结果**——没有重复,也不会发生意外数据丢失。 - **Bronze 层仅支持追加。** 每次 API 拉取都会写入一个带有新时间戳的文件 (`fuel_
标签:Databricks, 供应链管理, 商业智能, 开源数据, 数据工程