Apache Kafka Triggers

AwaitMessageTrigger

The AwaitMessageTrigger is a trigger that will consume messages polled from a Kafka topic and process them with a provided callable. If the callable returns any data, a TriggerEvent is raised.

For parameter definitions take a look at AwaitMessageTrigger.

KafkaMessageQueueTrigger

The KafkaMessageQueueTrigger is a dedicated interface class for Kafka message queues that extends the common MessageQueueTrigger. It is designed to work with the KafkaMessageQueueProvider and provides a more specific interface for Kafka message queue operations while leveraging the unified messaging framework.

For parameter definitions take a look at KafkaMessageQueueTrigger

For how to use the trigger, refer to the documentation of the Apache Kafka Message Queue Trigger

KafkaSharedStreamTrigger

The KafkaSharedStreamTrigger lets an AssetWatcher watch Kafka topics. Triggers with the same topics and kafka_config_id share one Kafka consumer in the triggerer instead of each opening their own. KafkaSharedStreamTrigger needs Airflow 3.3 or later.

The Kafka connection named by kafka_config_id must set enable.auto.commit to false in its extra field, or the trigger refuses to start. The shared consumer commits the offset of a message only after the trigger events from that message are saved to the metadata database. With auto-commit on, the Kafka client would commit offsets on a timer, possibly before the trigger events are saved. A triggerer crash at that point would lose the message, because Kafka would not deliver it again.

The same message can still reach a trigger more than once, for example when the triggerer restarts after saving the trigger events but before committing the offset. In that case the watched asset gets more than one asset event for the message.

from airflow.providers.apache.kafka.triggers.shared_stream import KafkaSharedStreamTrigger
from airflow.sdk import DAG, Asset, AssetWatcher

trigger = KafkaSharedStreamTrigger(topics=["orders"], kafka_config_id="kafka_default")
asset = Asset("kafka_orders", watchers=[AssetWatcher(name="orders_watcher", trigger=trigger)])

with DAG(dag_id="process_orders", schedule=[asset]):
    ...

For parameter definitions take a look at KafkaSharedStreamTrigger.

Was this entry helpful?