-
Notifications
You must be signed in to change notification settings - Fork 0
Pipeline Module ZH
通用的 Beam 轉換流程。這裡完全不知道 session、order 或任何其他商業概念——每一個掛勾點都是呼叫端提供的函式。
PipelineSchema(不可變 dataclass):metrics_namespace、allowed_stages、allowed_statuses。會被傳遞到每個需要知道管線遙測詞彙的階段;build_telemetry_event 會拒絕任何不在這些封閉集合內的 stage/status,讓打錯字的情況立即失敗而非默默流入下游。
parsing.py — PARSE_NORMALIZE 階段
ParseNormalize(PTransform)包裝 ParseNormalizeDoFn:解碼 bytes/text、json.loads,再執行呼叫端的 normalize_fn。標籤:parsed(主輸出)、unparseable(只有原因代碼——絕不含原始位元組)。若傳入 schema,還會產生 stage_timing/execution_log。
validation.py — VALIDATE 階段
Validate(PTransform)包裝 ValidateDoFn:呼叫端的 validate_fn(event) -> (is_valid, reason_code),依紀錄種類累加合法/不合法計數器,標籤為 valid(主輸出)/invalid。只有結構性/範圍檢查該放在 validate_fn 裡——任何需要分組後才看得到的檢查(例如重複偵測)應放進 build_outcome_fn。
aggregation.py — AGGREGATE 階段
AggregateAndEmit(PTransform):KeyBy(呼叫端的 key_fn)→ GroupByKey → BuildOutcome,呼叫呼叫端的 build_outcome_fn(key, records),並預期回傳 {"metrics": ..., "telemetry": [...], "invalid_summary": ...}。每個非 None 的輸出在被標記並產出前,都會先經過 safety.leak_guard.assert_safe 檢查(標籤:metrics/telemetry/invalid_summary)。tag_origin(event, origin) 是個小工具,用來標記紀錄來源(例如 "valid"/"excluded"),讓 build_outcome_fn 能在同一組內分辨不同紀錄。
main.py — 組裝邏輯
-
build_pipeline(...)— 針對常見情境的便利包裝:parse → validate → 標記 valid/invalid → key → group →build_outcome_fn,中間沒有額外的過濾階段。回傳一組具名 PCollection 的 dict(metrics、telemetry、invalid_summary、unparseable,若enable_logging=True則再加上execution_log/stage_timing)。 -
run_local(...)— 在DirectRunner下執行pipeline_builder,讀取本地 JSONL 檔案,並將每個具名輸出各自寫入一個本地 JSONL 檔案。完全不會碰觸任何雲端 API。 -
run_cloud(...)— 提交一個真實的、從 Pub/Sub 讀取的串流DataflowRunner任務;呼叫端必須已先通過自己的雲端關卡檢查。回傳{"job_id", "job_name"}。 -
apply_streaming_window(...)—GlobalWindows+ 一次性的AfterProcessingTime觸發器,ACCUMULATING模式,讓GroupByKey能對「每次執行一個有界批次」的串流來源正確觸發。
若管線需要在驗證與彙總之間加入額外的商業規則過濾階段(例如訂單有效性政策),應直接組合 ParseNormalize/Validate/AggregateAndEmit,而非呼叫 build_pipeline——參見 ecommerce/pipeline.py。
-
write_jsonl/read_jsonl— 單純的本地 JSONL 讀寫,永遠安全。 -
build_cloud_pubsub_sink/build_cloud_bigquery_sink— 各自接受一個無參數的gate_checkcallable,除非雲端關卡開啟,否則必須拋出例外;且必須在建立任何 client 之前呼叫。BigQuery sink 為 append-only、CREATE_NEVER(資料表必須已存在)。
BaseToolkitPipelineOptions 註冊本地模式參數(--input_jsonl、--output_dir、--run_id)、雲端模式參數(--input_subscription、--telemetry_topic、--metrics_table、--run_events_table——皆無預設值),以及雲端確認關卡參數(--enable_cloud_mode、--execute_cloud_pipeline、--confirm_cloud_execution)。
serialize_metric_row/decimal_to_json_string — 將 Decimal 欄位轉為決定性的字串以供 JSON 輸出。這裡沒有任何 metric 公式;那屬於呼叫端/領域邏輯(參見 ecommerce/metrics.py)。
build_telemetry_event(...) — 建立一列消毒過的遙測紀錄,依 schema 的封閉集合驗證 stage/status,再送入 assert_safe。
build_run_manifest(...) — 為一次本地或雲端執行建立消毒過的摘要 dict(run_id、mode、runner、時間、狀態、數量,以及任意 **extra)。