Wetask 的下一個更新將帶來持久且具圍欄的外部工作者、原子批次結算,以及第一方 Go、TypeScript 和 Python 用戶端。

Wetask 最初的目標很簡單:將通常圍繞背景工作系統所組合的元件,整合進單一 Go 語言的執行環境中。

這表示任務佇列、排程器、快取、共識、操作 API 以及可觀測性,而不需要在每次部署時分別使用獨立的 broker、結果後端以及排程服務。

Wetask 的下一個更新將把這個概念延伸到嵌入式處理器之外。應用程式將能夠在另一個程序、容器、語言或執行環境中執行自己的工作者,而 Wetask 仍負責管理持久的任務狀態。

Bring your own workers

當執行無法或不應在 Wetask 伺服器內部進行時,外部工作者會很有用:

  • 執行機器學習工作的 Python 工作者
  • 呼叫應用程式服務的 TypeScript 工作者
  • 處理高吞吐量後端工作的 Go 工作者
  • GPU 或記憶體專用工作者池
  • 獨立部署且具獨立發佈週期的工作者
  • 與佇列及控制平面程序隔離的工作者

第一方 Go、TypeScript 及 Python 用戶端提供工作者 API,而 HTTP 及 gRPC 仍可供自訂整合使用。

工作者可以依佇列與任務類型請求批次、更新長時間執行的作業、回報成功,或將失敗分類為暫時性、永久性或已取消。憑證可以限定於特定佇列與任務類型。

Fencing matters more than another Ack endpoint

至少一次傳遞的系統必須假設任務可能再次送達。工作者可能暫停、失去連線、超過租約,然後在另一個工作者已收到相同任務後繼續執行。

Wetask 為每次認領的嘗試提供一次性圍欄令牌。只有當任務、認領、令牌、主體、嘗試次數與協調者仍然相符時,才會接受完成。一旦嘗試過期或被取代,舊的工作者便無法完成它。

令牌本身只會回傳給工作者。Wetask 會儲存其 SHA-256 雜湊值。

這無法讓外部副作用變成恰好一次。如果工作者在確認任務前對卡片進行扣款並當機,仍可能導致重試。處理程序必須使用任務 ID 與嘗試次數,讓下游作業具有冪等性。

圍欄所能保證的是,過時的工作者無法覆寫目前嘗試的狀態。

Durable and idempotent settlement

成功與失敗的結算會在 Wetask 回傳收據前被持久化。重試相同的確認會回傳原始收據,而不會重複完成或重試該任務。

失敗處理是任務感知的:

  • transient:當仍有重試次數時,安排另一次嘗試
  • permanent:結束該任務並將其記錄至死信佇列
  • cancelled:結束執行而不重試
  • 租約到期時會自動套用已設定的重試與 DLQ 政策

工作者可以原子地確認或拒絕完整的批次。Wetask 會先驗證每一道圍欄。如果其中一項已過時或無效,該批次中所有未結算的項目都不會被變更。

批次結算也能透過一次 WAL 同步記錄整批,減少持久化儲存的負擔。

Benchmarks

我們對完整的處理中生命週期進行基準測試:

  1. 提交任務
  2. 以外部工作者認領任務
  3. 確認任務

該基準測試並未執行應用程式工作,也不包含 HTTP 或 gRPC 的網路負擔。它僅衡量任務執行環境本身。

環境:

  • Apple M3
  • macOS darwin/arm64
  • 回報 8 路平行處理的 Go 基準測試程序
  • 批次大小為 1、8 與 32
  • 每種情境三個樣本
  • 三秒基準測試視窗

指令:

go test ./internal/task \
  -run '^$' \
  -bench '^BenchmarkWorkerClaimAck($|Parallel$)' \
  -benchmem \
  -benchtime=3s \
  -count=3

Enter fullscreen mode Exit fullscreen mode

In-memory lifecycle

下表顯示每秒任務數的中位數:

情境 早期基準 目前結果 差異
Batch 1 74,933 77,123 +2.9%
Batch 8 84,434 86,393 +2.3%
Batch 32 84,661 86,783 +2.5%
Parallel, single-task operations 64,429 63,443 -1.5%

Batch 32 完成完整生命週期的時間約為每任務 11.52 microseconds

記憶體內的提升是有意保持適度。沒有磁碟障礙需要攤提,而且每項任務仍需 ID、認領令牌、狀態轉換、佇列記錄、結果保留、指標以及生命週期事件。平行處理的結果基本上持平,因此保留而未隱藏。

Durable lifecycle

持久模式包含已接受提交、認領與結算的 WAL 同步:

情境 中位數吞吐量 相對於早期 86.34 tasks/s 基準
Batch 1 80.08 tasks/s 0.93x
Batch 8 645.7 tasks/s 7.5x
Batch 32 2,402 tasks/s 27.8x
Parallel, single-task operations 136.4 tasks/s 1.58x

在 Batch 32 時,持久處理每任務約需 416 microseconds

理論上的批次上限為 32x,因為相同的持久化障礙由 32 個任務共用。實測的 27.8x 提升達到該上限的約 87%。剩餘時間用於 WAL 序列化、任務狀態、指標以及磁碟同步本身。

這些是 Wetask 自身的量測結果,並非宣稱 Wetask 比所有 broker 更快。公平的競爭性基準測試必須使用相同的酬載、傳輸層、批次大小、持久化政策、複製以及確認語意。

How Wetask compares

RabbitMQ quorum queues

RabbitMQ quorum queues 是成熟且以 Raft 複製的佇列,專為資料安全與領導者故障轉移而設計。發佈者確認與手動消費者確認提供穩健的生產基礎。

RabbitMQ 在複製佇列持久性、操作歷史、協定支援以及生態系規模上仍領先 Wetask。

Wetask 的目標是提供更針對任務的體驗:結果、重試、期限、取消、死信、排程、快取功能,以及工作者圍欄,皆可透過單一執行環境與 API 取得。

參考資料:
RabbitMQ quorum queues

NATS JetStream

JetStream 是最接近的協定層比較對象。它提供拉取與推送消費者、明確的 Ack、Nak、延遲 Nak、Term、進行中確認、AckWait、backoff 以及最大傳遞次數。

JetStream 在叢集串流儲存、發布/訂閱、複製以及高吞吐量訊息傳遞方面明顯更成熟。

Wetask 的差異在於嚴格的嘗試圍欄。JetStream 的文件指出,晚到的確認可能在訊息已被重新傳遞給另一個訂閱者後仍被接受。Wetask 則會拒絕來自已被取代的嘗試的結算。

參考資料:
NATS JetStream consumers

AWS SQS and Google Cloud Pub/Sub

受管理的佇列移除了大部分 broker 操作,並提供彈性、多可用區的基礎架構。

SQS 使用可見性逾時與收據控柄。標準佇列仍屬至少一次傳遞,消費者必須容忍重複傳遞。

Google Cloud Pub/Sub 的恰好一次拉取訂閱提供特別有趣的比較:重新傳遞後,只有最新的確認 ID 仍然有效。這在概念上接近 Wetask 的圍欄模型。

Wetask 適合需要自行託管、可攜性、本地部署以及任務特定狀態的團隊。當受管理的服務與供應商操作的高可用性更為重要時,雲端佇列是更好的選擇。

參考資料:

Apache Kafka

Kafka 是一個有序、可重播的事件日誌。它更適合事件串流、大量保留的歷史記錄、分割區排序以及許多獨立的消費者群組。

傳統的 Kafka 消費者提交分割區偏移量,而不是獨立地結算任意工作。單一記錄的重試可能影響分割區進度,而任務結果、每任務租約以及死信行為通常需要應用程式慣例或額外的主題。

Wetask 是背景工作的更直接模型。Kafka 則是持久事件串流的更強大模型。

參考資料:
Kafka delivery semantics

Celery with RabbitMQ

Celery 仍是 Python 應用程式的成熟選擇。它擁有龐大的生態系、熟悉的任務裝飾器、重試、Canvas 工作流,以及多年的操作知識。

其典型的生產架構結合 Celery 工作者、如 RabbitMQ 這樣的 broker、結果後端,以及 Celery Beat 作為排程器。

Wetask 的主張是更小的操作面:單一 Go 執行環境提供佇列、任務狀態、排程器、快取、API 以及管理,而外部工作者仍可使用 Python、TypeScript 或 Go 撰寫。

參考資料:
Celery documentation

Temporal

Temporal 解決的問題比任務佇列更大。它專為持久的多步驟工作流而設計,可在程序故障與基礎架構中斷期間恢復。

Temporal 是長時間執行的商業流程、saga、持久計時器以及工作流歷史的更佳選擇。Wetask 則是針對一般非同步工作、排程工作以及應用程式快取而刻意設計的更簡潔方案。

參考資料:
Temporal documentation

Where Wetask fits

Wetask 並非試圖取代所有訊息系統。

當日誌即產品時選擇 Kafka;當持久工作流編排即產品時選擇 Temporal;當外包操作是優先事項時選擇受管理的雲端佇列;當成熟的複製訊息基礎比整合式應用程式執行環境更重要時選擇 RabbitMQ 或 JetStream。

Wetask 適合想要以下特性的團隊:

  • 緊湊的自託管任務平台
  • 持久的背景工作
  • 伺服器程序之外的工作者
  • 嚴格的過時工作者拒絕
  • 任務結果、重試、期限以及 DLQ 管理
  • 無需額外服務堆疊即可使用的排程與快取功能
  • HTTP、gRPC、Go、TypeScript 以及 Python 整合

Availability

外部工作者已規劃於下一個 Wetask alpha 版本中推出。第一個版本將聚焦於持久的單節點使用以及受控的生產試點。多節點協調已存在,但任務佇列 WAL 擁有權目前為每節點;複製的任務擁有權以及永久節點遺失恢復仍屬後續強化工作。

此界線是刻意設定的。alpha 版本的目標是將完整的工作者體驗放入實際專案中,發布可重現的證據,並讓生產回饋塑造下一階段。

如果您正在評估背景工作基礎架構,最有用的回饋並非「哪個基準測試數字最大?」,而是:「哪些失敗語意、操作模型以及工作者體驗符合您實際需要執行的系統?」