이벤트 기반 스케줄링

이벤트 기반 스케줄링 (Event-driven scheduling)

이 페이지는 버전 3.0에 추가된 이벤트 기반 스케줄링을 다뤄요. 외부 이벤트 기반으로 DAG를 트리거할 수 있게 해주는 이벤트 기반 스케줄링은, 미리 정의된 시간 기반 스케줄이 아니라 실시간 데이터 변경·메시지·시스템 신호에 반응해야 하는 현대 데이터 아키텍처에서 특히 유용해요. AssetWatcherBaseEventTrigger 기반의 호환 트리거 작성법과 사용 사례를 설명해요.

출처: 문서

본문

버전 3.0에 추가됨.

Apache Airflow는 이벤트 기반 스케줄링을 허용하며, 미리 정의된 시간 기반 스케줄이 아니라 외부 이벤트에 기반해 DAG가 트리거되게 해요. 워크플로가 실시간 데이터 변경, 메시지, 시스템 신호에 반응해야 하는 현대 데이터 아키텍처에서 특히 유용해요.

Asset-Aware Scheduling에서 설명한 대로 asset을 사용해, 특정 외부 이벤트가 발생할 때 DAG가 실행을 시작하도록 구성할 수 있어요. Asset은 외부 이벤트와 DAG 실행 사이의 의존성을 수립하는 매커니즘을 제공해, 워크플로가 외부 환경의 변화에 동적으로 반응하도록 보장해요.

AssetWatcher 클래스는 이 매커니즘에서 중요한 역할을 해요. 메시지 큐 같은 외부 이벤트 소스를 모니터링하고, 관련 이벤트가 발생하면 asset 업데이트를 트리거해요. Asset 정의의 watchers 파라미터는 하나의 asset에 여러 AssetWatcher 인스턴스를 연결해, 다양한 이벤트 소스에 반응할 수 있게 해요.

자세한 내용과 예제는 common.messaging provider 문서를 참고해요.

이벤트 기반 스케줄링이 지원하는 트리거

BaseTrigger에서 상속하는 모든 triggers가 이벤트 기반 스케줄링에 사용될 수 있는 것은 아니에요. BaseTrigger에서 상속하는 모든 트리거와 달리, BaseEventTrigger에서 상속하는 부분집합만 호환돼요. 이 제한의 이유는 일부 트리거가 이벤트 기반 스케줄링용으로 설계되지 않았고, 그것을 DAG 스케줄링에 사용하면 의도하지 않은 결과를 초래할 수 있기 때문이에요.

BaseEventTrigger는 스케줄링에 사용되는 트리거가 이벤트 기반 패러다임을 준수하도록 보장해서, 예상치 못한 DAG 동작 없이 외부 이벤트 변경에 적절히 반응하게 해요.

이벤트 기반 호환 트리거 작성

트리거를 이벤트 기반 스케줄링과 호환되게 하려면 BaseEventTrigger에서 상속해야 해요. 이 맥락에서 트리거로 작업하는 세 가지 주요 시나리오가 있어요:

  1. 새 이벤트 기반 트리거 만들기: 지원되지 않는 이벤트 소스에 새 트리거가 필요하다면, BaseEventTrigger에서 상속하는 새 클래스를 만들고 그 로직을 구현해야 해요.
  2. 기존 호환 트리거 적응하기: 기존 트리거(BaseTrigger에서 상속한)가 이벤트 기반 스케줄링과 이미 호환된다는 것이 입증되면, 기본 클래스를 BaseTrigger에서 BaseEventTrigger로 바꾸기만 하면 돼요.
  3. 기존 비호환 트리거 적응하기: 기존 트리거가 이벤트 기반 스케줄링과 호환되지 않는 것 같으면, 새 트리거를 만들어야 해요. 이 새 트리거는 BaseEventTrigger에서 상속하고 이벤트 기반 스케줄링과 올바르게 작동하도록 해야 해요. 두 트리거가 일부 공통 코드를 공유한다면 기존 트리거에서 상속할 수도 있어요.

형제 트리거 간 하나의 폴 공유

버전 3.3에 추가됨.

서로 다른 asset의 여러 AssetWatcher 인스턴스가 같은 upstream 리소스를 읽는 트리거를 지원할 때 — 플래그 파일의 디렉터리, 폴링 REST 엔드포인트, 그런 멱등성·구독자 부작용 소스 — triggerer는 그렇지 않으면 트리거마다 하나의 독립 폴 루프를 띄워요. 스무 명의 구독자를 가진 공유 소스라면 스무 개의 폴 루프, 스무 개의 연결, 주기당 스무 세트의 API 호출을 의미해요. 정확한 범위는 아래 "적합한 upstream"을 참고해요.

BaseEventTrigger는 형제 트리거가 하나의 기반 폴을 공유하면서 각 트리거가 자신의 DB 행, run_trigger task, per-instance 필터링을 유지하는 opt-in 경로를 지원해요. 참여하려면 서브클래스가 세 개의 훅을 오버라이드해요:

  • shared_stream_key() — 공유 upstream을 식별하는 키를 반환해요(보통 문자열 튜플). 키가 같게 비교되는 트리거는 하나의 폴을 공유해요. None 반환(기본값)은 opt-out이에요 — 트리거가 이전처럼 정확히 자체 독립 run() 루프를 실행해요. 반환값은 triggerer가 이 트리거를 시작할 때 한 번 읽어요. 수명 중간에 바꿔도 그룹 구성원에는 효과가 없으므로, 폴을 공유해야 하는 형제들은 처음부터 같은 키를 반환해야 해요. 키는 결정적이어야 해요 — time.time()이나 uuid.uuid4() 같은 per-call 값이 아니라 설정 필드에서 파생해야 해요. 그룹의 수명 전반에 걸쳐 비교가 안정적이어야 하니까요.
  • open_shared_stream() — triggerer가 shared-stream 그룹당 한 번 구동해 upstream에서 raw 이벤트를 생성하는 @classmethod 코루틴. Triggerer가 하나의 트리거 kwargs를 재사용해 공유 폴을 구동하므로, shared_stream_key에 참여하는 필드 값에만 의존해요.
  • filter_shared_stream() — 브로드캐스트된 raw 스트림을 소비하고 이 트리거가 발생시킬 TriggerEvent 인스턴스를 생성하는 인스턴스 메서드. Per-trigger 필터링(예: 이 인스턴스의 filename과 일치하는 이벤트만)이 여기 있어요.

예제: 공유 inbox 디렉터리에 per-asset 플래그 파일이 나타날 때 발생하는 DirectoryFileDeleteTrigger:

from collections.abc import AsyncIterator, Hashable
from typing import Any

from airflow.triggers.base import BaseEventTrigger, TriggerEvent

class DirectoryFileDeleteTrigger(BaseEventTrigger):
    def __init__(self, *, directory, filename, poke_interval=5.0):
        super().__init__()
        self.directory = directory
        self.filename = filename
        self.poke_interval = poke_interval

    def shared_stream_key(self) -> Hashable | None:
        # All triggers on the same directory + cadence share one scan.
        return ("directory-scan", self.directory, self.poke_interval)

    @classmethod
    async def open_shared_stream(cls, kwargs: dict[str, Any]) -> AsyncIterator[Any]:
        # Drives one directory listing loop per group.
        ...

    async def filter_shared_stream(self, shared_stream: AsyncIterator[Any]) -> AsyncIterator[TriggerEvent]:
        # Each instance fires only for its own filename.
        async for snapshot in shared_stream:
            if self.filename in snapshot["names"]:
                yield TriggerEvent(...)
                return

이 트리거를 사용한 완전한 예제는 airflow.example_dags.example_asset_with_watchers에 있으며, 두 형제 DirectoryFileDeleteTrigger watcher가 같은 DAG 안의 독립형 FileDeleteTrigger watcher와 함께 하나의 디렉터리 스캔을 공유해요.

무엇이 공유되고 무엇이 안 되는가

공유는 이름이 암시하는 것보다 더 좁아요:

  • 공유됨 (shared_stream_key당 하나): open_shared_stream async generator와 그 upstream I/O — 예를 들어 디렉터리의 실제 iterdir 호출이나 폴링 REST API 호출.
  • 공유되지 않음 (트리거당 하나): Trigger DB 행, 트리거 인스턴스, run_trigger asyncio task, filter_shared_stream async generator. 각 AssetWatcher는 여전히 UI와 메타데이터 데이터베이스에서 자신의 트리거로 나타나요.

즉 절약은 폴 루프·upstream I/O 계층에서 오는 것이지, 영속·스케줄링 계층이 아니에요.

적합한 upstream

공유 스트림 패턴에 잘 맞는 것들:

  • 멱등성/읽기 전용 읽기 — 디렉터리 스캔, 폴링 REST API.
  • 구독자 부작용 정리 — 트리거의 per-event 동작(unlink, 로컬 표시, …)이 공유 producer 핸들과 독립적으로 subscriber가 소유한 API를 통해 이루어지는 경우.
  • 메시지 브로커 upstream(Kafka, SQS, Pub/Sub, Azure Service Bus) — producer가 모든 subscriber가 메시지를 처리한 후에만 commit/delete/ack해야 하는 경우 — 아래 설명하는 ack 채널을 사용해요.

Producer측 ack 채널

producer가 모든 subscriber가 이벤트를 처리한 후에만 진행(commit, delete, ack)해야 하는 upstream의 경우, create_shared_stream_producer()를 오버라이드해 SharedStreamProducer를 반환해요. 이 팩토리가 오버라이드되면 manager는 ack 모드에 들어가요. Subscriber 쪽은 변하지 않아요: filter_shared_stream은 fast path와 정확히 같은 raw 이벤트를 받으므로 같은 필터 코드가 두 모드 모두에서 작동해요 — 프레임워크는 각 subscriber의 소비 진행과 그것이 파생한 trigger event의 영속에서 브로커가 언제 진행할 수 있는지 추론해요.

  1. Manager는 shared-stream 그룹당 한 번 create_shared_stream_producer(kwargs)를 호출해요. 반환된 producer는 한 번의 폴 수명 동안 브로커 연결을 소유해요. open_stream 안에서 연결을 지연해서 열어요 — 팩토리에서가 아니라요.
  2. Producer의 open_stream(event, broker_payload) 튜플을 생성해요. 여기서 broker_payload는 producer가 나중에 필요한 것(예: SQS receipt handle, Kafka offset, Pub/Sub ack ID)이에요.
  3. 각 subscriber의 filter_shared_stream은 fast path에서와 정확히 같은 raw 이벤트를 받아요. Subscriber는 이벤트를 지나친 다음(다음 raw 이벤트를 당기거나 구독 해제) 그리고 그것이 이벤트에서 파생한 모든 TriggerEvent가 메타데이터 데이터베이스에 영속된 후에 이벤트를 해결돼요(resolve).
  4. fan-out 세트의 모든 subscriber가 이벤트를 해결했을 때(지나침, 구독 해제, 타임아웃, 큐 오버플로로), manager는 이벤트의 lane에서 완전히 해결된 이벤트의 연속 접두사를 가지고 await producer.advance(batch)를 호출해요 — offset을 commit하고, SQS 메시지를 삭제하는 등. 배치의 각 항목은 이벤트의 broker_payload와 subscriber들이 어떻게 해결했는지에 대한 per-event 카운트를 가진 AdvanceOutcome를 실은 AdvanceItem이에요.

이벤트 거부: subscriber의 필터는 그 이벤트를 처리하는 동안 reject_shared_stream_event()를 호출해, raw 이벤트에서 trigger event를 생성하는 대신 적극적으로 거부할 수 있어요. 이는 비자발적 실패와 구별돼요. failed 카운트는 subscriber가 제때 끝나지 않았거나(ack 타임아웃) 뒤처졌음을(큐 오버플로) 의미하며, 올바른 대응은 보통 재전달이에요. rejected 카운트는 subscriber가 이벤트가 trigger event를 생성하면 안 되고 최종적으로 폐기돼야 한다고 결정했음을 의미하며, 올바른 대응은 재전달이 아니라 dead-letter하거나 nack하는 것이에요. AdvanceOutcome은 두 카운트를 별도로 보고하므로 producer가 advance에서 올바른 per-broker 작업을 적용할 수 있어요:

  • Azure Service Bus: rejected가 0이 아니면 메시지를 dead-letter 하고, failed만 0이 아니면 버리고(브로커가 재전달), 모든 subscriber가 이벤트를 수락했으면 완료해요.
  • Pub/Sub: reject 시 메시지를 nack하고, 그렇지 않으면 ack해요.

프레임워크는 카운트만 보고해요 — 절대 스스로 dead-letter, nack, 재전달하지 않아요. 그 브로커 특정 결정은 전적으로 producer의 advance에 있어요.

reject_shared_stream_event는 필터가 ack 모드에서 raw 이벤트를 처리하는 동안에만 의미가 있어요(그 이벤트의 바인딩 창이 열려 있을 때). Fast path, 독립 run(), 또는 두 raw 이벤트 사이에서 호출되면 경고를 기록하고 아무것도 하지 않아요. 브로커 진행에 영향을 줄 것이 없으니까요. 이벤트를 즉시 해결하므로 영속할 것이 없고, reject는 영속 게이트를 기다리지 않아요.

is_clean은 브로드캐스트 시 온라인이었던 모든 subscriber가 이벤트를 수락했을 때만 True예요: reject 없음, 실패 없음, 그리고 적어도 하나의 subscriber가 ack했음. 단일 reject나 실패는 False로 만들고, subscriber가 0인 브로드캐스트(all-zero 카운트)도 그렇죠 — 아무도 수락하지 않았으므로 producer가 커밋할 것이 없어요.

예제 — 필터 안에서 reject:

from airflow.triggers.shared_stream import reject_shared_stream_event

async def filter_shared_stream(self, shared_stream):
    async for raw in shared_stream:
        if raw.get("malformed"):
            # Never produce a trigger event from this; have the broker
            # dead-letter it rather than redeliver it.
            reject_shared_stream_event()
            continue
        yield TriggerEvent(raw)

순서 보장: 기본적으로 모든 이벤트는 같은 lane에 속해요. 배치의 항목은 이벤트 순서이며 그 lane의 연속 해결 접두사를 이뤄요. lane 안에서 배치는 엄격히 순서대로 도착해요 — 다음 advance는 이전 호출이 반환된 후에만 await돼요. advance가 예외를 발생시키면 에러가 기록되고 전체 shared-stream 그룹이 종료돼요: 모든 subscriber가 실패 센티널을 받고, 브로커가 절대 커밋되지 않은 offset에서 재전달해요(조용한 데이터 누락보다는 크고 안전한 실패). 일시적 실패에서 더 우아하게 복구하려면 producer가 어느 offset이 다시 커밋해도 안전한지 추적해야 해요. Producer는 get_advance_lane을 오버라이드해 그 순서를 lane 안으로 좁힐 수 있어요: lane 값이 같게 비교되는 이벤트는 서로 상대적으로 이벤트 순서로 배치·진행되고, 다른 lane의 이벤트는 서로 기다리지 않아요. 어느 쪽이든 한 번에 최대 하나의 advance 호출만 await되며, Kafka offset commit 같은 누적 방식은 배치의 마지막 항목만 커밋하면 돼요. 폴이 끝나면 manager가 await producer.aclose()를 한 번 던져요(best-effort).

예제 — SQS 스타일 producer:

from collections.abc import AsyncIterator, Sequence
from typing import Any

from airflow.triggers.base import BaseEventTrigger, TriggerEvent
from airflow.triggers.shared_stream import AdvanceItem, SharedStreamProducer

class SqsSharedStreamProducer(SharedStreamProducer):
    def __init__(self, queue_url: str):
        self.queue_url = queue_url
        self.client = None

    async def open_stream(self) -> AsyncIterator[tuple[Any, Any]]:
        # Open the connection here, not in the trigger's factory.
        self.client = await create_sqs_client()
        while True:
            messages = await poll_sqs(self.client, self.queue_url)
            for msg in messages:
                yield msg["Body"], msg["ReceiptHandle"]

    async def advance(self, batch: Sequence[AdvanceItem]) -> None:
        # Called with one lane's batch of fully resolved messages.
        for receipt_handle, outcome in batch:
            if outcome.is_clean:
                await delete_sqs_message(self.client, self.queue_url, receipt_handle)
            # Otherwise leave the message for the visibility timeout to redeliver.

    async def aclose(self) -> None:
        if self.client is not None:
            await self.client.close()

class SqsSharedTrigger(BaseEventTrigger):
    def __init__(self, *, queue_url: str, region: str | None = None):
        super().__init__()
        self.queue_url = queue_url
        self.region = region

    def serialize(self):
        return (
            f"{type(self).__module__}.{type(self).__qualname__}",
            {"queue_url": self.queue_url, "region": self.region},
        )

    def shared_stream_key(self):
        return ("sqs", self.queue_url)

    @classmethod
    def create_shared_stream_producer(cls, kwargs) -> SqsSharedStreamProducer:
        return SqsSharedStreamProducer(kwargs["queue_url"])

    async def filter_shared_stream(self, shared_stream):
        async for raw in shared_stream:
            if self.region is None or raw.get("region") == self.region:
                yield TriggerEvent(raw)

    async def run(self):
        yield TriggerEvent({})

예제 — 파티션 간 Kafka 누적 커밋. Kafka 커밋은 파티션 내에서 커밋된 offset까지의 모든 offset을 인정하므로, 이전 이벤트가 아직 pending인 동안 같은 파티션의 이후 이벤트가 커밋되는 것은 안전하지 않아요 — 다른 파티션의 이벤트는 문제가 되지 않아요. get_advance_lane에서 (topic, partition)을 반환하면 순서 보장을 정확히 그 세분성으로 좁혀요: 각 파티션의 커밋은 순서대로 유지되고, 느린 파티션이 더 이상 다른 파티션의 커밋을 지연시키지 않아요:

class KafkaSharedStreamProducer(SharedStreamProducer):
    def __init__(self, topics: list[str]):
        self.topics = topics
        self.consumer = None

    async def open_stream(self):
        # Auto-commit must be off (the Kafka default is on), or the
        # consumer commits on its own schedule and the ack channel
        # no longer controls what the broker considers delivered.
        self.consumer = await create_kafka_consumer(self.topics, enable_auto_commit=False)
        async for message in self.consumer:
            yield message.value, (message.topic, message.partition, message.offset)

    def get_advance_lane(self, broker_payload):
        topic, partition, _offset = broker_payload
        return topic, partition

    async def advance(self, batch):
        # The batch is one lane's — here, one partition's — contiguous
        # resolved prefix, in event order, so committing the offset of
        # its last item covers the whole batch and can never skip past
        # an event that is still pending.
        #
        # If this method raises, the whole shared-stream group is
        # terminated and the broker redelivers from the last committed
        # offset — no silent data skip.
        #
        # Inspect each item's outcome before committing: non-clean events
        # (rejects, failures, or a zero-subscriber broadcast) should go to
        # a dead-letter queue rather than be committed as delivered.
        to_dlq = []
        for broker_payload, outcome in batch:
            if not outcome.is_clean:
                to_dlq.append(broker_payload)
        if to_dlq:
            handle_dlq(to_dlq)
        topic, partition, offset = batch[-1].broker_payload
        await self.consumer.commit(topic, partition, offset + 1)

    async def aclose(self):
        if self.consumer is not None:
            await self.consumer.stop()

필터 작성자에게 한 가지 제약은 raw 이벤트와 그것에서 파생된 trigger event 사이의 바인딩이에요: raw 이벤트에서 파생된 모든 TriggerEvent를, 공유 스트림에서 다음 raw 이벤트를 당기기 전에 생성해야 해요 — 단순한 필터 루프가 어차피 하는 일이에요.

Snapshot-at-fan-out: 주어진 이벤트를 해결해야 하는 subscriber의 집합은 이벤트가 브로드캐스트되는 순간 고정돼요. 이벤트가 디스패치된 후 합류하는 subscriber는 그 이벤트의 pending 세트에 추가되지 않아요.

Per-event ack 타임아웃: subscriber가 ack 타임아웃(기본 5분, [triggerer] shared_stream_ack_timeout 설정 옵션으로 구성 가능) 안에 이벤트 처리를 끝내지 못하면 — 여전히 그 이벤트에 있거나, 그로부터 파생된 일부 trigger event가 영속 확인되지 않았거나 — manager가 그 subscriber의 트리거를 강제 실패시켜요. 다른 subscriber는 영향을 받지 않아요. 그들이 해결하면 producer는 정상 진행해요. Ack 타임아웃은 manager 레벨 안전망이며, 네이티브 브로커 세션이나 visibility 타임아웃을 대체하지 않아요. Subscriber 관점에서 강제 실패는 filter_shared_stream 안의 shared_stream 반복자가 일으키는 AckTimeout(airflow.triggers.shared_stream에서 import 가능)으로 나타나요. 그대로 전파시켜도 괜찮아요 — 트리거는 표준 트리거 실패 경로로 실패해요. subscriber가 실패하기 전에 정리를 실행해야 할 때만 잡아요.

Triggerer 재시작: 해결 상태는 메모리에만 있어요. Triggerer 재시작 후 브로커는 진행되지 않은 메시지를 재전달해요. 따라서 subscriber는 멱등성이어야 해요.

같은 키를 공유하는 여러 트리거가 함께 재시작하면, 가장 먼저 재구독하는 트리거가 새 그룹을 만들고 폴링이 즉시 시작돼요. 나중에 재구독하는 트리거는 늦은 subscriber로 합류해서(이미 브로드캐스트된 이벤트의 snapshot 밖), 첫 번째 구독과 자신 사이의 창에 커밋된 이벤트를 놓칠 수 있어요. [triggerer] shared_stream_cohort_grace_period를 양수 초(예: 2.0)로 설정하면 새 그룹이 만들어진 후 폴링 시작을 지연시켜, 어떤 이벤트도 브로드캐스트되기 전에 동시 재구독이 합류할 시간을 줘요. 이것은 best-effort 창이에요 — 느리게 재합류하는 트리거가 이벤트를 놓칠 위험을 줄이지만 없애지는 않아요.

내구성: 브로커 진행은 영속에 게이트가 걸려요. Subscriber의 해결은 이벤트에서 파생한 모든 TriggerEvent가 메타데이터 데이터베이스에 저장된 후에만 완료돼요. 확인은 다음 상태 동기화 시 트리거 러너에 도달하는데, 보통 1~2초 안에요. 확인이 절대 도착하지 않으면 — triggerer가 크래시했거나, 이벤트를 영속할 수 없었거나 — ack 타임아웃이 이벤트를 실패시키고, producer는 커밋하지 않으며, 브로커가 재전달해요. 따라서 실패는 중복 전달을 일으킬 수 있지만 이벤트 손실은 절대 없어요. 멱등한 subscriber가 중복을 흡수해요. 같은 절충이 이벤트가 아직 확인을 기다리는 동안 그룹이 멈출 때 적용돼요(예: 마지막 subscriber가 이벤트를 생성한 직후 구독 해제): pending 진행은 버려지고 브로커가 그 이벤트를 재전달해요.

ack 모드의 shared_stream_subscriber_queue_size: 설정 경계는 여전히 subscriber당 처리되지 않은 raw 이벤트를 지배해요. Manager는 다음 upstream 이벤트를 당기기 전에 미해결 해결을 기다리지 않아요; back-pressure는 큐에 바인딩돼요 — 큐가 꽉 찬 subscriber는 강제 실패돼요. 큐는 주로 subscriber의 필터가 실행되기 전의 버스트 전달을 막아요. 반면 브로커 진행은 각 lane 안에서 순서대로 디스패치돼요: pending subscriber가 있는 이벤트는 같은 lane의 모든 이후 이벤트의 진행을 지연시키고 — 뒤의 해결된 이벤트는 다음 배치로 누적 — ack 타임아웃이 그 대기를 경계로 해요.

공유가 활성인지 검증하기

Triggerer는 각 shared-stream 그룹의 생성을 기록하고, poll task 이름을 그 키로 짓습니다:

Shared stream group started key=('directory-scan', '/tmp/region-flags', 5.0)
asyncio task name: shared-stream-poll[('directory-scan', '/tmp/region-flags', 5.0)]

공유가 활성이라면, 얼마나 많은 subscriber가 합류하는지와 무관하게 고유 키당 정확히 하나의 Shared stream group started 줄이 보여야 해요. 대신 subscriber당 하나의 로그 줄이 보인다면 키가 같게 비교되지 않는 것일 수 있어요 — 형제들 사이에서 shared_stream_key가 같은 값을 반환하는지 확인해요.

느린 subscriber 오버플로

Shared-stream 그룹의 각 subscriber는 경계가 있는 인메모리 큐를 가져요. 폴 루프가 subscriber의 filter_shared_stream이 소비할 수 있는 것보다 빠르게 이벤트를 생성하면, 큐가 차고 그 트리거가 _SubscriberOverflow로 실패해요 — 무한 메모리 증가보다 의도된 fail-fast예요.

Subscriber가 반복적으로 오버플로하면 두 가지 방법으로 해결할 수 있어요:

  • [triggerer] shared_stream_subscriber_queue_size를 올려 오버플로 임계값에 도달하기 전에 필터에 더 많은 여유를 줘요.
  • shared_stream_key()를 재설계해 한 그룹을 공유하는 형제 트리거를 줄여요 — 더 좁은 그룹은 subscriber 하나가 소비해야 하는 이벤트 비율을 줄여요.

둘 다 producer 처리량과 per-subscriber 소비율 사이의 불일치를 줄여요.

무한 스케줄링 피하기

일부 트리거가 이벤트 기반 스케줄링과 호환되지 않는 이유는, 외부 리소스가 주어진 상태에 도달하기를 기다리기 때문이에요. 예시:

  • 저장소 서비스에 파일이 존재하기를 기다리기
  • Job이 성공 상태가 되기를 기다리기
  • 데이터베이스에 행이 존재하기를 기다리기

이런 조건에서 스케줄링하면 무한 재스케줄링으로 이어질 수 있어요. 조건이 참이 되면 오랫동안 참으로 유지될 가능성이 높기 때문이에요.

예를 들어 특정 job이 "success" 상태에 도달할 때 실행되도록 스케줄링된 DAG를 생각해 보세요. Job이 성공하면 보통 그 상태를 유지해요. 결과적으로 triggerer가 조건을 확인할 때마다 DAG가 반복적으로 트리거돼요.

또 다른 예는 S3 버킷에 특정 파일이 존재하는지 확인하는 S3KeyTrigger예요. 파일이 생성되면 "파일 X가 버킷 Y에 존재하나?"라는 조건이 참으로 유지되므로 트리거가 매 확인마다 계속 성공해요. 이는 트리거 매커니즘이 실행될 때마다 DAG가 무기한 트리거되게 해요.

커스텀 트리거를 만들 때, 한번 충족되면 영구적으로 참으로 유지되는 조건을 사용하는 데 주의하세요. 이는 의도치 않게 무한 DAG 실행을 초래하고 시스템을 압도할 수 있어요.

이벤트 기반 DAG의 사용 사례

  • 데이터 수집 파이프라인: 저장 시스템에 새 데이터가 도착하면 ETL 워크플로를 트리거해요.
  • 머신러닝 워크플로: 새 데이터셋이 사용 가능해지면 모델 학습을 시작해요.
  • IoT 및 실시간 분석: 센서 데이터, 로그, 애플리케이션 이벤트에 실시간으로 반응해요.
  • 마이크로서비스·이벤트 기반 아키텍처: 서비스 간 메시지에 기반해 워크플로를 오케스트레이션해요.

더 알아보기 (Learn more)