tests.system.amazon.aws.example_kinesis_message_queue

Example Dag demonstrating event-driven scheduling with Amazon Kinesis Data Streams.

NOTE: This file serves as an example Dag and reference for AssetWatcher configuration. It is NOT an automated end-to-end integration test: running this file directly (e.g. via pytest system-test harness) initiates a manual DagRun where triggering_asset_events is empty. Validating the complete event-driven chain requires a running Airflow triggerer, an active Kinesis stream, and external records producing AssetEvents.

Pre-requisites: 1. An active Amazon Kinesis Data Stream must exist and be accessible by the configured AWS connection. 2. The Airflow triggerer must be running with the common.messaging provider installed. 3. This is an event-driven Dag triggered by an AssetWatcher; it does not produce records to itself.

Attributes

STREAM_NAME

AWS_CONN_ID

AWS_REGION

trigger

kinesis_asset

test_run

Functions

process_kinesis_records(**context)

Process and decode incoming records triggered from the Amazon Kinesis stream.

Module Contents

tests.system.amazon.aws.example_kinesis_message_queue.STREAM_NAME[source]
tests.system.amazon.aws.example_kinesis_message_queue.AWS_CONN_ID[source]
tests.system.amazon.aws.example_kinesis_message_queue.AWS_REGION[source]
tests.system.amazon.aws.example_kinesis_message_queue.trigger[source]
tests.system.amazon.aws.example_kinesis_message_queue.kinesis_asset[source]
tests.system.amazon.aws.example_kinesis_message_queue.process_kinesis_records(**context)[source]

Process and decode incoming records triggered from the Amazon Kinesis stream.

tests.system.amazon.aws.example_kinesis_message_queue.test_run[source]

Was this entry helpful?