活用 Pub/Sub 与 CDC:在 PostgreSQL 中订阅表变更以替代 ETL 数据同步
2026/9/4 1:24:20 网站建设 项目流程

在資料工程實務中,最讓團隊疲憊的往往不是複雜的報表 SQL,而是大量“把一張表原樣搬到另一個儲存”的同步任務。這類任務通常叫 ETL,如果更精確一點,它主要是 ODS 層的搬表工作。Tabsdata 提出的 “Pub/Sub for Tables to Replace ETL Pipelines” 之所以值得關注,不是因為又要取代哪個調度工具,而是在提醒資料團隊重新思考一個問題:能不能不要把資料表當作批次拷貝的對象,而是把它當作一種可以被訂閱的即時事件流?

這個思路的價值,在於把“資料表變更”變成和訊息佇列裡的 event 一樣自然。下游不需要頻繁跑全量對比,也不需要依賴每一張表都有可靠的 updated_at。只要能訂閱一張資料表的主題,插入、更新、刪除都能被即時接收,並以冪等方式寫入目標儲存。

下面會先分析傳統批次 ETL 和 ODS 層為什麼會造成維護成本,再說明表格級 Pub/Sub 的事件模型,然後用一個 PostgreSQL 搭配 Python 的最小案例跑通“訂閱一張表”的完整過程,最後補上生產環境落地時的參數選擇、排錯路徑和最佳實踐。

1. 把 ETL 拆開看待:真正被吐槽的是“搬表”,不是“計算”

很多團隊一提到 ETL,就以為是在處理數據轉換邏輯。但在日常開發裡,佔用資料工程師時間最多的,往往是 ODS 層的同步作業:把來源庫的 user、order、product 等表,按定時任務複製到倉庫或湖裡。

麻煩的是,這樣的同步任務看起來簡單,實際維護成本卻集中在幾個容易被忽略的細節上。

1.1 ODS 層本質上是“定時搬表”的苦力作業

ODS 層在傳統數倉裡通常是第一站。它的目標很樸素:先把原始資料完整搬進倉,再留給後續 DWD、ADS 層做清洗和彙總。但真正落地時,很多 ODS job 會寫成這樣:

-- 常見的 ODS 增量同步 SQL,按 updated_at 抓取變化資料 CREATE TEMP TABLE tmp_user AS SELECT id, name, email, status, updated_at FROM source_db.user WHERE updated_at >= %(last_sync_time)s; BEGIN; DELETE FROM ods_user WHERE id IN (SELECT id FROM tmp_user); INSERT INTO ods_user SELECT id, name, email, status, updated_at FROM tmp_user; COMMIT;

這段語句看起來很直接,但它預設了三個前提:

  • 來源表一定存在且可靠的updated_at
  • 來源表的資料只會按照時間往後變化。
  • 下游可以接受“先刪後插”的同步方式。

任何一個前提不成立,這個 job 都會出問題。更麻煩的是,問題常常不會立即暴露,而是要等到下游分析出結果異常才被發現。

1.2 只靠 updated_at 做增量,究竟丟失了哪些信息

假設有一張 order 表,業務方手動修正了一筆歷史訂單,但沒有更新 updated_at。增量同步就會漏掉這筆資料。實際項目中,這種“應用層沒寫時間戳”的情況遠比想像中常見。

即使應用層每次都正確更新 updated_at,仍然會有幾個問題:

  • 物理刪除無法被增量查詢捕獲。如果來源表直接DELETE FROM order WHERE id = 123,同步 job 可能永遠不知道這筆資料消失了。
  • 同一時間點內修改多筆資料時,按 updated_at 篩選容易出現邊界重複或漏數。
  • 當表比較大時,每隔幾分鐘就做一次WHERE updated_at >= time掃描,會給來源資料庫帶來明顯壓力。
  • 業務表 DDL 一調整,例如新增欄位、修改型別,同步 SQL 就可能報錯,需要人工介入。

所以很多 ODS 層的維護工作,其實是在不斷“打補丁”:加刪除標記、加版本號、加 start_time/end_time,最後把本來很簡單的搬表任務改造成一套複雜的自研同步系統。

1.3 哪些 ETL 應該被替換,哪些不應該被替換

不能說 Pub/Sub for Tables 要替代所有 ETL。正確的判斷方式是區分工作類型。

用途傳統批次 ETL表格級 Pub/Sub
大規模離線彙總適合,例如 T+1 報表不適合,仍需要批次計算
高頻熱點表同步延遲高,容易產生業務髒資料適合,交易提交後即可推送
支援多個下游即時消費同一張表每個下游各跑一套同步邏輯適合,一份變更流多個消費者
需要回放某個時間點的資料狀態重跑成本高透過訂閱位移和快照機制處理
大量明細資料一次性灌入倉庫可以批次初始化初始化完成後建議切到增量流

理解這個邊界後,才能真正看懂 Tabsdata 的定位:它想解決的是“表同步”這件事,而不是取代資料分析場景裡的複雜轉換。

2. 表格級 Pub/Sub 的核心:把資料表變成可以訂閱的 Topic

Pub/Sub 這個概念在微服務架構中很常見。一個服務發布訂單建立事件,另一個服務訂閱這個事件,雙方不需要直接呼叫 API。把這個模型搬運到資料表上,就誕生了“表即 Topic”的設計。

Tabsdata 所說的 Pub/Sub for Tables,核心是讓一個資料表本身變成發布源。當資料表的資料發生新增、修改、刪除時,系統把這個變更封裝成一個事件,發布到對應的佇列或主題上。任何註冊過訂閱的消費者,都能收到這個變更。

2.1 一張資料表的事件模型不是“一條訊息”,而是一個行變更

普通訊息佇列裡的 event 往往代表一次業務行為,例如“使用者點擊了按鈕”。但資料表事件表達的是“某一行的狀態變了”。

如果來源表是這樣一張使用者表:

CREATE TABLE app_user ( id BIGINT PRIMARY KEY, name VARCHAR(64) NOT NULL, email VARCHAR(128) NOT NULL, status VARCHAR(16) NOT NULL DEFAULT 'active', updated_at TIMESTAMPTZ NOT NULL DEFAULT now() );

當應用執行UPDATE app_user SET status = 'blocked' WHERE id = 1001時,表格級 Pub/Sub 發布的事件不會是“資料更新成功”這種含糊訊息,而是一筆帶有 before 和 after 的行變更事件。

換成常見 CDC 事件格式,大概長這樣:

{ "op": "u", "source": { "db": "shop_db", "table": "public.app_user", "ts_ms": 1700000000000 }, "before": { "id": 1001, "name": "Alice", "email": "alice@example.com", "status": "active" }, "after": { "id": 1001, "name": "Alice", "email": "alice@example.com", "status": "blocked" } }

這裡的op通常取值為:

  • c表示新增。
  • u表示更新。
  • d表示刪除。
  • r表示初始化快照時讀取出來的資料。

事件裡必須包含主鍵或唯一鍵,因為消費者要靠它決定要把哪一筆資料合併到目標表。

2.2 為什麼資料表事件不只是“資料庫觸發器”

第一次接觸這個概念的人可能會想:資料表變更,用資料庫 trigger 寫入一張 event 表不就行了?

這背後有一個很重要的差異:觸發器是應用層邏輯,它只解決“能感知變更”,卻無法覆蓋資料庫底層的複製語義。例如:

  • 觸發器依賴表結構,如果 DDL 改動頻繁,函數很容易失效。
  • 觸發器沒有統一的 schema 演進機制。
  • 大流量的表如果每個操作都觸發一次函數,PostgreSQL 的效能開銷會上升。
  • 如果消費者掛了,trigger 也不具備下游位移回放能力。

因此,真正的表格級 Pub/Sub 通常會基於資料庫日誌,也就是 CDC。PostgreSQL 用 WAL,MySQL 用 binlog,SQL Server 用 transaction log。讀取資料庫操作日誌的好處是:可以在事務提交後以近乎原始的方式取得變更,不影響來源庫的業務行為,也不需要每一次同步都用 SQL 去掃來源表。

2.3 表變更流和普通 Topics 的差異

如果是訂閱一個普通業務 Topic,消費者通常只需要處理“事件發生了一次”。但訂閱一張資料表時,消費者需要理解:接收到的不是業務語義,而是資料庫行狀態的變遷。

因此消費端邏輯更像“做即時複製”:

  • op = c代表目標表應該插入一行。
  • op = u代表目標表應該更新某一行。
  • op = d代表目標表應該刪除某一行。
  • op = r通常代表初始化快照,目標表也應該按照新行處理。

這種設計讓表 Pub/Sub 相較於批次 ETL 多了一個關鍵能力:任何新增消費者,只需要重播 Topic 中的歷史變更,就能把目標表還原到當前狀態。不需要回到來源資料庫再做一次全量掃描關聯。

3. 一個最小案例:把 PostgreSQL 使用者表發佈到 Topic

概念聽完後,需要跑一個最小案例,才能真正體會“表變更被即時訂閱”的過程。

這裡用 PostgreSQL 的LISTEN/NOTIFY做簡化示範。它只是為了理解流程,不是生產環境的替代方案。真正生產環境會使用 Debezium、Flink CDC、Maxwell 或雲端資料庫 CDC 服務。先跑通最小案例後,再換到底層機制會容易得多。

環境準備:

  • PostgreSQL 14 以上

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询