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.