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¶
Functions¶
|
Process and decode incoming records triggered from the Amazon Kinesis stream. |