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 sqs as 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, LATEST only sees records that arrive after the watcher starts polling. If the watcher is down or new shards are discovered, earlier records may be skipped. Use TRIM_HORIZON to 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.

tests/system/amazon/aws/example_kinesis_message_queue.py[source]

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

Was this entry helpful?