Apache Kafka Message Queue¶
Apache Kafka Queue Provider¶
Implemented by KafkaMessageQueueProvider
The Apache Kafka Queue Provider is a BaseMessageQueueProvider that uses
Apache Kafka as the underlying message queue system.
It allows you to send and receive messages using Kafka topics in your Airflow workflows with MessageQueueTrigger common message queue interface.
It uses
kafkaas scheme for identifying Kafka queues.For parameter definitions take a look at
AwaitMessageTrigger.from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger from airflow.sdk import Asset, AssetWatcher trigger = MessageQueueTrigger( scheme="kafka", # Additional Kafka AwaitMessageTrigger parameters as needed topics=["my_topic"], apply_function="module.apply_function", bootstrap_servers="localhost:9092", ) asset = Asset("kafka_queue_asset", watchers=[AssetWatcher(name="kafka_watcher", trigger=trigger)])For a complete example, see:
tests.system.common.messaging.kafka_message_queue_trigger
Apache Kafka Message Queue Trigger¶
Implemented by KafkaMessageQueueTrigger
Inherited from MessageQueueTrigger
Wait for a message in a queue¶
Below is an example of how you can configure an Airflow Dag to be triggered by a message in Apache Kafka.
from airflow.providers.apache.kafka.triggers.msg_queue import KafkaMessageQueueTrigger
from airflow.providers.standard.operators.empty import EmptyOperator
from airflow.sdk import DAG, Asset, AssetWatcher
def apply_function(message):
val = json.loads(message.value())
print(f"Value in message is {val}")
return True
# Define a trigger that listens to an Apache Kafka message queue
trigger = KafkaMessageQueueTrigger(
topics=["test"],
apply_function="example_dag_kafka_message_queue_trigger.apply_function",
kafka_config_id="kafka_default",
apply_function_args=None,
apply_function_kwargs=None,
poll_timeout=1,
poll_interval=5,
)
# Define an asset that watches for messages on the queue
asset = Asset("kafka_queue_asset_1", watchers=[AssetWatcher(name="kafka_watcher_1", trigger=trigger)])
with DAG(dag_id="example_kafka_watcher_1", schedule=[asset]) as dag:
EmptyOperator(task_id="task")
How it works¶
Kafka Message Queue Trigger: The
KafkaMessageQueueTriggerlistens for messages from Apache Kafka Topic(s).
2. Asset and Watcher: The Asset abstracts the external entity, the Kafka queue in this example.
The AssetWatcher associate a trigger with a name. This name helps you identify which trigger is associated to which
asset.
3. Event-Driven Dag: Instead of running on a fixed schedule, the Dag executes when the asset receives an update (e.g., a new message in the queue).
For how to use the trigger, refer to the documentation of the Messaging Trigger
The apply_function¶
The apply_function is applied to every message polled from the Kafka topic(s). If it returns a truthy
value, that value is used as the payload of the TriggerEvent. Otherwise, the trigger keeps polling.
It is required when using the Kafka queue provider, while the Kafka sensors also accept None, in which
case the raw message value (decoded as UTF-8) is used as the event payload.
The function must be passed as a string in Python dot notation, for example
"my_package.my_module.my_function", not as the function object, because trigger arguments are
serialized to the metadata database. The function is imported in the triggerer at runtime, so the module
must be importable there, and changes to the function require a triggerer restart.
Additional arguments can be passed using apply_function_args and apply_function_kwargs. The message
is always passed as the last positional argument:
# my_package/my_module.py
import json
from confluent_kafka import Message
def my_function(prefix: str, message: Message, threshold: int = 0) -> str | None:
val = json.loads(message.value())
if val["amount"] > threshold:
return f"{prefix}{val}"
# In your Dag file
from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger
trigger = MessageQueueTrigger(
scheme="kafka",
topics=["my_topic"],
apply_function="my_package.my_module.my_function",
apply_function_args=["received:"],
apply_function_kwargs={"threshold": 100},
)