-
Notifications
You must be signed in to change notification settings - Fork 0
Overview ZH
beam-pipeline-toolkit 是一個獨立、可重複使用的 Apache Beam 管線框架。這個 repo 只包含通用的建構模組——沒有 demo 應用程式、沒有儀表板、也沒有特定公司的商業邏輯。
| 套件 | 職責 |
|---|---|
pipeline/ |
通用 Beam 轉換流程:parse → validate → aggregate → sink,以及本地/雲端執行器的組裝邏輯 |
controlplane/ |
Dataflow 控制平面 API:提交任務、請求取消任務、讀取任務證據(中繼資料 + counters) |
log/ |
選用的執行日誌(消毒過的例外事件)+ 階段耗時記錄 + 唯讀的 Cloud Logging 讀取器 |
safety/ |
遞迴式 PII/機密外洩防護,以及「碰觸雲端資源前必須確認」的關卡機制 |
ecommerce/ |
選用、可直接改造使用的電商領域層,建構於上述四個套件之上 |
大多數事件管線專案都會出現同樣的 Beam 轉換流程(parse → validate → aggregate → sink)與同樣的 Dataflow submit/stop/evidence 機制。真正不同的是商業邏輯:什麼樣的紀錄算合法、紀錄如何分組、一列 metric 代表什麼意義。這個框架把那個邊界精確地推給呼叫端提供的函式(normalize_fn、validate_fn、key_fn、build_outcome_fn)與少數幾個小型設定物件(PipelineSchema),因此當商業規則改變時,框架本身完全不需要修改。
ecommerce/ 就是最好的證明:它是一個完整、可運作的管線,使用的正是外部呼叫端會用到的同一套公開 API——pipeline/、controlplane/、log/、safety/ 內完全不知道 session、order 或 customer 的存在。
原始 JSON 訊息
│
▼
ParseNormalize (pipeline/parsing.py) -- JSON 解碼 + normalize_fn
│ 標籤:parsed | unparseable
▼
Validate (pipeline/validation.py) -- 呼叫端提供的 validate_fn
│ 標籤:valid | invalid
▼
AggregateAndEmit (pipeline/aggregation.py) -- key_fn、GroupByKey、build_outcome_fn
│ 標籤:metrics | telemetry | invalid_summary
▼
sinks (pipeline/sinks.py) -- 本地 JSONL,或有關卡保護的 BigQuery/Pub/Sub
當傳入 PipelineSchema 且 enable_logging=True 時,以上每個階段都可以選擇性地產生 stage_timing 與 execution_log 事件(參見日誌模組)。
python -m venv .venv
.venv\Scripts\activate
pip install -r requirements.txt
pytest
完整可執行範例(本地執行、雲端提交、讀取任務證據)請見 beam_pipeline_toolkit/README.md。
固定於 requirements.txt:apache-beam[gcp]、google-auth、requests、pytest、build。完整的間接相依版本則凍結於 requirements.lock.txt。