Skip to content

Overview ZH

northrails edited this page Jul 21, 2026 · 1 revision

English version

總覽

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_fnvalidate_fnkey_fnbuild_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

當傳入 PipelineSchemaenable_logging=True 時,以上每個階段都可以選擇性地產生 stage_timingexecution_log 事件(參見日誌模組)。

快速開始

python -m venv .venv
.venv\Scripts\activate
pip install -r requirements.txt
pytest

完整可執行範例(本地執行、雲端提交、讀取任務證據)請見 beam_pipeline_toolkit/README.md

相依套件

固定於 requirements.txtapache-beam[gcp]google-authrequestspytestbuild。完整的間接相依版本則凍結於 requirements.lock.txt

Clone this wiki locally