Byte Ebi's Logo

Byte Ebi 🍤

每天一小口,蝦米變鯨魚

使用 S3KeySensor 讓 Airflow 不用等一輩子檔案

只要能等到檔案生成,我什麼都願意做

Ray

Airflow 裡「等外部非同步產生的檔案才能繼續」是很常見的情境
但「多久都願意等」,對負載來說是多麼沉重的一句話

既有問題

目前經手一個系統整合流程,透過 Airflow 排程,主要步驟:

  1. 接收參數
  2. 呼叫外部 API,產生一份要交換的文件
  3. 產生完成後,繼續往下進行後續流程

其中「等檔案產生完成」這個環節,前後在不同專案裡用過兩種做法

第一版

直接呼叫外部 API 等對方把檔案產生完才回應

檔案產生要花一段時間,常常等到 API timeout 就被切斷連線,整個 task 直接失敗
當時的解法是把 timeout 設定拉長,治標不治本

第二版

改成外部服務產生檔案的同時,把檔案生成完成的狀態寫進 Redis
當初的工程師手刻了一個 sensor 等 Redis key 出現才往下走,藉此繞過 timeout 問題

看似解決了 timeout,但實際運行時踩了幾個坑:

  • worker 所在的 ECS instance 使用率過低時會被自動回收,上面可能還有 task 在等待
  • Redis key 會過期,沒有即時檢查到就會查不到狀態,即使檔案其實已經產生完成

上面兩個問題還會量子糾纏,變成更麻煩的:

ECS instance 被回收之後沒人繼續檢查,key 就這樣默默過期了

直接變成薛丁格的狀態,就算失敗也不一定是真的失敗

目前的應對方式是拉長 scale-incooldown period 並加大預設 worker 數量
一樣是治標不治本

解決方案

S3KeySensor + deferrable 來拯救大家啦!

在 Airflow 裡「等外部非同步產生的檔案才能繼續」這種情境其實很常見
內建的 S3KeySensor 就是設計來處理這件事的!

可以直接偵測「S3 上某個 key 是否存在
而且是用 S3 object 作為偵測目標,所以不會有 Redis key 過期的問題

加上 S3KeySensor 是 Airflow 官方 provider 維護的標準實作
經過 Airflow 官方驗證使用 deferrable=True 時不會佔用 worker slot
也不會有 worker 因為 ECS instance 被回收而中斷的問題

支援三種模式:

  • mode="poke"預設值):task 會佔用 worker slot 直到條件成立為止
  • mode="reschedule":檢查時會重新排程一個 worker 執行,在檢查中間會釋放 worker slot
  • deferrable=True(預設是 False):把「等待」交給獨立的 triggerer process 處理,不佔用 worker

三者之間是有優先順序的,可以查看 S3KeySensor.execute() 的原始碼

apache/airflow GitHub - providers/amazon/…/sensors/s3.py

def execute(self, context: Context) -> None:
    if not self.deferrable:
        super().execute(context)   # mode="poke"/"reschedule" 只在這裡起作用
    else:
        if not self.poke(context=context):
            self._defer()          # 交給 triggerer 背景等待,不佔用 worker

deferrable=True 時只會 poke 一次,沒滿足就直接丟給 triggerer,無論 mode 設什麼都不會生效
mode 只有在 deferrable=False 時才有意義

本篇的範例選擇了 deferrable=True 就不需要、也不應該再另外設定 mode

其他相關參數的預設值:

  • poke_interval:預設 60 秒,每次檢查間隔
  • timeout:預設 604800 秒(7 天),逾時仍等不到就判定失敗

概念驗證

寫了一個簡單的 DAG 驗證整個機制:

import logging
from datetime import UTC, datetime, timedelta

from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor
from airflow.sdk import DAG, Param, Variable, get_current_context, task
from airflow.sdk.definitions.param import ParamsDict
from services.aws.s3_client import S3Client

logger = logging.getLogger(__name__)

default_args = {
    "owner": "Developer",
    "depends_on_past": False,
    "start_date": datetime(2026, 9, 22, tzinfo=UTC),
    "retry_delay": timedelta(minutes=1),
}

with DAG(
    dag_id="s3_key_sensor_test",
    default_args=default_args,
    schedule=None,
    catchup=False,
    max_active_runs=1,
    dag_display_name="S3KeySensor 概念驗證",
    tags=["AWS"],
    doc_md="""\
        ## 測試 `S3KeySensor` 用法

        1. `get_s3_path`:產生要監控的 S3 路徑
        2. `wait_for_key`:等待 `get_s3_path` 產生的 S3 object 出現才會往下走
        3. `fetch_value`:偵測到檔案後用既有的 `S3Client` 讀取內容,回傳長度 (bytes)
        4. `print_value`:接收 `fetch_value` 的回傳值並印出來
    """,
    params=ParamsDict(
        {
            "bucket_name": Param(
                default="",
                type="string",
                title="S3 Bucket",
                description="S3 bucket name",
            ),
            "bucket_key": Param(
                default="",
                type="string",
                title="Bucket Key",
                description="S3 object key",
            ),
        }
    ),
) as dag:

    @task(task_id="get_s3_path", multiple_outputs=True)
    def get_s3_path() -> dict[str, str]:
        """模擬正式流程中路徑是由程式邏輯算出來的,這裡先直接用 Trigger DAG 帶入的 params 代替

        multiple_outputs=True 讓每個 dict key 分別存成獨立的 XCom 鍵值,下游才能用 s3_path["bucket_name"] 這種下標語法取值,否則會拿到 None

        Returns:
            dict[str, str]: 各自對應一個 XCom key
        """
        params = get_current_context().get("params", {})
        return {
            "bucket_name": params["bucket_name"],
            "bucket_key": params["bucket_key"],
        }

    s3_path = get_s3_path()

    wait_for_object = S3KeySensor(
        task_id="wait_for_object",
        bucket_name=s3_path["bucket_name"],  # type: ignore[index]
        bucket_key=s3_path["bucket_key"],  # type: ignore[index]
        deferrable=True,  # 是否交給 triggerer 執行,不佔用 worker slot
        poke_interval=30,  # 每次檢查間隔秒數,正式建議 >= 60
        timeout=3600,  # 逾時秒數
    )

    @task(task_id="fetch_object", retries=0)
    def fetch_object(bucket_name: str, bucket_key: str) -> int:
        content = S3Client(bucket_name=bucket_name).get_object(bucket_key)

        logger.info("[S3KeySensor] 在 S3 %s/%s 偵測到檔案", bucket_name, bucket_key)
        return len(content)

    @task(task_id="print_object_length", retries=0)
    def print_value(content_length: int) -> None:
        print(f"[S3KeySensor] 成功取得 {content_length} bytes")

    fetch_object_task = fetch_object(
        s3_path["bucket_name"],  # type: ignore[index]
        s3_path["bucket_key"],  # type: ignore[index]
    )

    wait_for_object.set_downstream(fetch_object_task)
    print_value(fetch_object_task)  # type: ignore[arg-type]

if __name__ == "__main__":
    dag.cli()

1. 取得 S3 路徑

正式流程中檔案上傳的路徑是由程式按照規則產生的
這邊為了方便測試,改成觸發時手動輸入 bucket_name/bucket_key

取得 S3 路徑

可以在 XCom 分頁看到 bucket_namebucket_key 各自被記錄下來
return_value 則是包成 dict 一起傳出去

這裡的關鍵是 multiple_outputs=True
它讓 dict 裡的每個 key 各自存成獨立的 XCom,下游才能用 s3_path["bucket_name"] 直接取值

沒有設定的話 XCom 就只會存一份 return_value
這時候 s3_path["bucket_name"] 拿到的其實是 None,而不是預期中的字串

2. 用 S3KeySensor 等待 S3 object 出現

透過 deferrable sensor 等待 get_s3_path 傳來的路徑真的出現檔案

deferred 狀態

從審計日誌可以清楚看到狀態變化:running -> deferred -> running -> success

一開始 task 執行後馬上轉成 deferred 狀態 中間等了將近 7 分鐘,這段期間完全不佔用 worker

3. 偵測到檔案後讀取內容

手動上傳檔案,模擬外部 API 生成檔案

上傳檔案到 S3

S3KeySensor 偵測到檔案出現後就從 deferred 恢復並往下執行

取得檔案內容

比對 log 和實際檔案內容長度一致

印出檔案長度

注意事項

使用 deferrable 模式需要環境有啟用 triggerer 元件,沒開的話 task 會卡在 deferred 永遠不會被喚醒

在 PoC 裡沒有另外設定 aws_conn_id,是直接吃 MWAA 執行角色的預設 boto3 憑證鏈

「worker 被回收就失敗」的問題不是 100% 消除,而是換成一個更輕量的依賴

During the deferred phase of execution, since work has been offloaded to the triggerer, the task no longer occupies a worker slot, and you have more free workload capacity.

等待邏輯改成跑在獨立的 triggerer process 上,所以變成依賴 triggerer 要活著
好消息是 deferred 狀態有存進 metadata DB,即使 triggerer process 重啟也不會直接判失敗,藉此達成高可用的目標

Airflow automatically re-schedules triggers that were on that host to run elsewhere.

當條件成立、trigger 觸發之後,也還是要交回一個 worker 才能真正執行完 task:

  • The trigger runs until it fires, at which point its source task is re-scheduled by the scheduler.
  • The scheduler queues the task to resume on a worker node.

所以說使用 deferrable 可以消除的是

  • 等待期間依賴的 worker process 隨時可能被 ECS 回收

換成風險低很多的

  • 依賴 triggerer 存活

    • deferred 狀態的任務存在 metadata DB
    • 就算其中一個 triggerer 掛掉 scheduler 也會自動把 trigger 轉給其他還活著的 triggerer
  • 恢復執行時仍需要排到一個 worker

    • 只在任務真正執行的時間內佔用 worker slot

參考來源:Deferrable Operators & Triggers — Apache Airflow 官方文件

最新文章

文章分類

Tag