Amazon Messaging Queues¶
Amazon SQS Queue Provider¶
Implemented by SqsMessageQueueProvider
The Amazon SQS Queue Provider is a BaseMessageQueueProvider that uses
Amazon Simple Queue Service (SQS) as the underlying message queue system.
It allows you to send and receive messages using SQS queues in your Airflow workflows with MessageQueueTrigger common message queue interface.
It uses
sqsas scheme for identifying SQS queues.For parameter definitions take a look at
SqsSensorTrigger.
from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger
from airflow.sdk import Asset, AssetWatcher
trigger = MessageQueueTrigger(
scheme="sqs",
# Additional AWS SqsSensorTrigger parameters as needed
sqs_queue="https://sqs.us-east-1.amazonaws.com/123456789012/my-queue",
aws_conn_id="aws_default",
)
asset = Asset("sqs_queue_asset", watchers=[AssetWatcher(name="sqs_watcher", trigger=trigger)])
For a complete example, see:
tests.system.amazon.aws.example_dag_sqs_message_queue_trigger
Amazon Kinesis Data Streams Provider¶
Implemented by KinesisMessageQueueProvider
The Amazon Kinesis Data Streams Provider is a BaseMessageQueueProvider that uses
Amazon Kinesis Data Streams as the underlying messaging system.
It enables event-driven scheduling with MessageQueueTrigger using scheme="kinesis".
For parameter definitions take a look at KinesisTrigger.
from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger
from airflow.sdk import Asset, AssetWatcher
trigger = MessageQueueTrigger(
scheme="kinesis",
stream_name="my-kinesis-stream",
aws_conn_id="aws_default",
)
watcher = AssetWatcher(name="kinesis_watcher", trigger=trigger)
asset = Asset("kinesis_stream_asset", watchers=[watcher])
Delivery semantics and considerations:
Record payload: Record data in the trigger event payload (
message_batch) is base64-encoded and must be decoded by consuming tasks.Shard iterator type: When no checkpoint exists,
LATESTonly sees records that arrive after the watcher starts polling. If the watcher is down or new shards are discovered, earlier records may be skipped. UseTRIM_HORIZONto process from the oldest available record.Checkpointing: Checkpointing shard progress is supported when a single asset is watched in an Airflow runtime providing an asset state store.
Best-effort delivery: Delivery is best-effort. In the event of triggerer restarts or transient failures, records may be re-delivered or missed around failure windows. It does not provide exactly-once guarantees.
Amazon Kinesis Data Streams Message Queue Trigger¶
Implemented by KinesisTrigger
Dispatched by MessageQueueTrigger for scheme="kinesis"
Wait for records in a stream¶
Below is an example of how you can configure an Airflow Dag to be triggered by records published to an Amazon Kinesis data stream.
trigger = MessageQueueTrigger(
scheme="kinesis",
stream_name=STREAM_NAME,
aws_conn_id=AWS_CONN_ID,
region_name=AWS_REGION,
shard_iterator_type="TRIM_HORIZON",
)
kinesis_asset = Asset(
f"kinesis://{STREAM_NAME}",
watchers=[AssetWatcher(name="kinesis_stream_watcher", trigger=trigger)],
)
@task
def process_kinesis_records(**context) -> None:
"""Process and decode incoming records triggered from the Amazon Kinesis stream."""
events = context["triggering_asset_events"].get(kinesis_asset, [])
if not events:
print(
"No triggering asset events found for this run. "
"When executed manually or via test runners without an active watcher, "
"no Kinesis records are delivered. In an event-driven environment, "
"the Airflow triggerer emits an AssetEvent containing Kinesis records."
)
return
for event in events:
message_batch = event.extra.get("payload", {}).get("message_batch", [])
for record in message_batch:
raw_data = base64.b64decode(record["Data"]).decode("utf-8")
print(
f"Received record: ShardId={record['ShardId']}, "
f"SequenceNumber={record['SequenceNumber']}, "
f"Data={raw_data}"
)
with DAG(
dag_id="example_kinesis_message_queue",
schedule=[kinesis_asset],
start_date=datetime(2025, 1, 1),
catchup=False,
tags=["example", "kinesis", "message_queue"],
) as dag:
process_kinesis_records()
For how to use the trigger, refer to the documentation of the Messaging Trigger