This is an automated email from the ASF dual-hosted git repository.

vincbeck pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new 66607ac5483 AIP-82: Add KinesisMessageQueueProvider (#73509)
66607ac5483 is described below

commit 66607ac5483cd50a4ffe00a15d92d085f6e47dc9
Author: Aaron Chen <[email protected]>
AuthorDate: Mon Sep 28 06:55:11 2026 -0700

    AIP-82: Add KinesisMessageQueueProvider (#73509)
---
 providers/amazon/docs/message-queues/index.rst     |  55 ++++++
 providers/amazon/provider.yaml                     |   1 +
 .../airflow/providers/amazon/aws/queues/kinesis.py |  48 ++++++
 .../airflow/providers/amazon/get_provider_info.py  |   5 +-
 .../amazon/aws/example_kinesis_message_queue.py    |  98 +++++++++++
 .../tests/unit/amazon/aws/queues/test_kinesis.py   | 185 +++++++++++++++++++++
 6 files changed, 391 insertions(+), 1 deletion(-)

diff --git a/providers/amazon/docs/message-queues/index.rst 
b/providers/amazon/docs/message-queues/index.rst
index 3d5af93769d..f752b7511ef 100644
--- a/providers/amazon/docs/message-queues/index.rst
+++ b/providers/amazon/docs/message-queues/index.rst
@@ -48,3 +48,58 @@ It allows you to send and receive messages using SQS queues 
in your Airflow work
 
 For a complete example, see:
 :mod:`tests.system.amazon.aws.example_dag_sqs_message_queue_trigger`
+
+
+Amazon Kinesis Data Streams Provider
+------------------------------------
+
+Implemented by 
:class:`~airflow.providers.amazon.aws.queues.kinesis.KinesisMessageQueueProvider`
+
+The Amazon Kinesis Data Streams Provider is a 
:class:`~airflow.providers.common.messaging.providers.base_provider.BaseMessageQueueProvider`
 that uses
+Amazon Kinesis Data Streams as the underlying messaging system.
+It enables event-driven scheduling with 
:class:`~airflow.providers.common.messaging.triggers.msg_queue.MessageQueueTrigger`
 using ``scheme="kinesis"``.
+For parameter definitions take a look at 
:class:`~airflow.providers.amazon.aws.triggers.kinesis.KinesisTrigger`.
+
+.. code-block:: python
+
+    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.
+
+.. _howto/triggers:KinesisMessageQueueTrigger:
+
+Amazon Kinesis Data Streams Message Queue Trigger
+-------------------------------------------------
+
+Implemented by 
:class:`~airflow.providers.amazon.aws.triggers.kinesis.KinesisTrigger`
+
+Dispatched by 
:class:`~airflow.providers.common.messaging.triggers.msg_queue.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.
+
+.. exampleinclude:: 
/../../amazon/tests/system/amazon/aws/example_kinesis_message_queue.py
+    :language: python
+    :start-after: [START howto_trigger_kinesis_message_queue]
+    :end-before: [END howto_trigger_kinesis_message_queue]
+
+For how to use the trigger, refer to the documentation of the
+:ref:`Messaging Trigger <howto/trigger:MessageQueueTrigger>`
diff --git a/providers/amazon/provider.yaml b/providers/amazon/provider.yaml
index 08425a38827..e9788a4fa81 100644
--- a/providers/amazon/provider.yaml
+++ b/providers/amazon/provider.yaml
@@ -1540,4 +1540,5 @@ cli:
   - airflow.providers.amazon.aws.cli.definition.get_aws_cli_commands
 
 queues:
+  - airflow.providers.amazon.aws.queues.kinesis.KinesisMessageQueueProvider
   - airflow.providers.amazon.aws.queues.sqs.SqsMessageQueueProvider
diff --git 
a/providers/amazon/src/airflow/providers/amazon/aws/queues/kinesis.py 
b/providers/amazon/src/airflow/providers/amazon/aws/queues/kinesis.py
new file mode 100644
index 00000000000..1a975b86360
--- /dev/null
+++ b/providers/amazon/src/airflow/providers/amazon/aws/queues/kinesis.py
@@ -0,0 +1,48 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+from __future__ import annotations
+
+from typing import TYPE_CHECKING
+
+from airflow.providers.amazon.aws.triggers.kinesis import KinesisTrigger
+from airflow.providers.common.compat.sdk import 
AirflowOptionalProviderFeatureException
+
+try:
+    from airflow.providers.common.messaging.providers.base_provider import 
BaseMessageQueueProvider
+except ImportError:
+    raise AirflowOptionalProviderFeatureException(
+        "This feature requires the 'common.messaging' provider to be installed 
in version >= 2.0.0. "
+        "Install it with: pip install 
'apache-airflow-providers-amazon[common.messaging]'"
+    )
+
+if TYPE_CHECKING:
+    from airflow.triggers.base import BaseEventTrigger
+
+
+class KinesisMessageQueueProvider(BaseMessageQueueProvider):
+    """
+    Configuration for Amazon Kinesis Data Streams integration with 
common-messaging.
+
+    Dispatches ``scheme="kinesis"`` to
+    :class:`~airflow.providers.amazon.aws.triggers.kinesis.KinesisTrigger`, 
which also defines
+    the accepted parameters.
+    """
+
+    scheme = "kinesis"
+
+    def trigger_class(self) -> type[BaseEventTrigger]:
+        return KinesisTrigger
diff --git a/providers/amazon/src/airflow/providers/amazon/get_provider_info.py 
b/providers/amazon/src/airflow/providers/amazon/get_provider_info.py
index 80e314ef7f1..388b4395e4d 100644
--- a/providers/amazon/src/airflow/providers/amazon/get_provider_info.py
+++ b/providers/amazon/src/airflow/providers/amazon/get_provider_info.py
@@ -1584,5 +1584,8 @@ def get_provider_info():
         ],
         "auth-managers": 
["airflow.providers.amazon.aws.auth_manager.aws_auth_manager.AwsAuthManager"],
         "cli": 
["airflow.providers.amazon.aws.cli.definition.get_aws_cli_commands"],
-        "queues": 
["airflow.providers.amazon.aws.queues.sqs.SqsMessageQueueProvider"],
+        "queues": [
+            
"airflow.providers.amazon.aws.queues.kinesis.KinesisMessageQueueProvider",
+            "airflow.providers.amazon.aws.queues.sqs.SqsMessageQueueProvider",
+        ],
     }
diff --git 
a/providers/amazon/tests/system/amazon/aws/example_kinesis_message_queue.py 
b/providers/amazon/tests/system/amazon/aws/example_kinesis_message_queue.py
new file mode 100644
index 00000000000..1852815ff56
--- /dev/null
+++ b/providers/amazon/tests/system/amazon/aws/example_kinesis_message_queue.py
@@ -0,0 +1,98 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+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.
+"""
+
+from __future__ import annotations
+
+import base64
+import os
+from datetime import datetime
+
+from airflow.providers.common.messaging.triggers.msg_queue import 
MessageQueueTrigger
+from airflow.sdk import DAG, Asset, AssetWatcher, task
+
+STREAM_NAME = os.getenv("KINESIS_STREAM_NAME", 
"airflow-kinesis-example-stream")
+AWS_CONN_ID = os.getenv("AWS_CONN_ID", "aws_default")
+AWS_REGION = os.getenv("AWS_REGION", "us-east-1")
+
+# [START howto_trigger_kinesis_message_queue]
+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()
+# [END howto_trigger_kinesis_message_queue]
+
+
+from tests_common.test_utils.system_tests import get_test_run  # noqa: E402
+
+# Needed to run the example DAG with pytest (see: 
contributing-docs/testing/system_tests.rst)
+test_run = get_test_run(dag)
diff --git a/providers/amazon/tests/unit/amazon/aws/queues/test_kinesis.py 
b/providers/amazon/tests/unit/amazon/aws/queues/test_kinesis.py
new file mode 100644
index 00000000000..a7e0c676848
--- /dev/null
+++ b/providers/amazon/tests/unit/amazon/aws/queues/test_kinesis.py
@@ -0,0 +1,185 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+from __future__ import annotations
+
+import importlib
+from unittest import mock
+
+import pytest
+
+from tests_common.test_utils.common_msg_queue import mark_common_msg_queue_test
+
+pytest.importorskip("airflow.providers.common.messaging.providers.base_provider")
+
+
+def test_message_kinesis_queue_create():
+    from airflow.providers.amazon.aws.queues.kinesis import 
KinesisMessageQueueProvider
+    from airflow.providers.common.messaging.providers.base_provider import 
BaseMessageQueueProvider
+
+    provider = KinesisMessageQueueProvider()
+    assert isinstance(provider, BaseMessageQueueProvider)
+    assert provider.scheme == "kinesis"
+
+
[email protected](
+    ("scheme", "expected_result"),
+    [
+        pytest.param("kinesis", True, id="kinesis_scheme"),
+        pytest.param("sqs", False, id="sqs_scheme"),
+        pytest.param("kafka", False, id="kafka_scheme"),
+        pytest.param("redis+pubsub", False, id="redis_scheme"),
+        pytest.param("unknown", False, id="unknown_scheme"),
+    ],
+)
+def test_message_kinesis_scheme_matches(scheme, expected_result):
+    from airflow.providers.amazon.aws.queues.kinesis import 
KinesisMessageQueueProvider
+
+    provider = KinesisMessageQueueProvider()
+    assert provider.scheme_matches(scheme) == expected_result
+
+
[email protected](
+    "queue",
+    [
+        pytest.param("kinesis://my-stream", id="kinesis_uri"),
+        
pytest.param("arn:aws:kinesis:us-east-1:123456789012:stream/my-stream", 
id="kinesis_arn"),
+        pytest.param("my-stream", id="stream_name_only"),
+    ],
+)
+def test_message_kinesis_queue_matches(queue):
+    from airflow.providers.amazon.aws.queues.kinesis import 
KinesisMessageQueueProvider
+
+    provider = KinesisMessageQueueProvider()
+    assert provider.queue_matches(queue) is False
+
+
+def test_message_kinesis_queue_trigger_class():
+    from airflow.providers.amazon.aws.queues.kinesis import 
KinesisMessageQueueProvider
+    from airflow.providers.amazon.aws.triggers.kinesis import KinesisTrigger
+
+    provider = KinesisMessageQueueProvider()
+    assert provider.trigger_class() == KinesisTrigger
+
+
+def test_message_kinesis_queue_trigger_kwargs():
+    from airflow.providers.amazon.aws.queues.kinesis import 
KinesisMessageQueueProvider
+
+    provider = KinesisMessageQueueProvider()
+    assert provider.trigger_kwargs("kinesis://my-stream", 
stream_name="my-stream") == {}
+
+
+def test_message_kinesis_missing_common_messaging_dependency():
+    import airflow.providers.amazon.aws.queues.kinesis as kinesis_mod
+    from airflow.providers.common.compat.sdk import 
AirflowOptionalProviderFeatureException
+
+    with mock.patch.dict("sys.modules", 
{"airflow.providers.common.messaging.providers.base_provider": None}):
+        with pytest.raises(
+            AirflowOptionalProviderFeatureException,
+            match=(
+                r"This feature requires the 'common\.messaging' provider to be 
installed in version "
+                r">= 2\.0\.0\. Install it with: pip install "
+                r"'apache-airflow-providers-amazon\[common\.messaging\]'"
+            ),
+        ):
+            importlib.reload(kinesis_mod)
+
+    importlib.reload(kinesis_mod)
+
+
+@mark_common_msg_queue_test
+class TestKinesisMessageQueueTriggerIntegration:
+    """Integration tests for KinesisMessageQueueProvider with 
MessageQueueTrigger and ProvidersManager."""
+
+    @pytest.mark.usefixtures("cleanup_providers_manager")
+    def test_provider_discovery(self):
+        from airflow.providers_manager import ProvidersManager
+
+        manager = ProvidersManager()
+        manager.initialize_providers_queues()
+        assert (
+            
"airflow.providers.amazon.aws.queues.kinesis.KinesisMessageQueueProvider"
+            in manager.queue_class_names
+        )
+
+    @pytest.mark.usefixtures("cleanup_providers_manager")
+    def test_message_queue_trigger_dispatch_to_kinesis(self):
+        from airflow.providers.amazon.aws.triggers.kinesis import 
KinesisTrigger
+        from airflow.providers.common.messaging.triggers.msg_queue import 
MessageQueueTrigger
+
+        trigger = MessageQueueTrigger(
+            scheme="kinesis",
+            stream_name="test-stream",
+            aws_conn_id="aws_test",
+            shard_iterator_type="TRIM_HORIZON",
+        )
+        assert isinstance(trigger.trigger, KinesisTrigger)
+        assert trigger.trigger.stream_name == "test-stream"
+        assert trigger.trigger.aws_conn_id == "aws_test"
+        assert trigger.trigger.shard_iterator_type == "TRIM_HORIZON"
+
+    @pytest.mark.usefixtures("cleanup_providers_manager")
+    def test_message_queue_trigger_serialize(self):
+        from airflow.providers.common.messaging.triggers.msg_queue import 
MessageQueueTrigger
+
+        trigger = MessageQueueTrigger(
+            scheme="kinesis",
+            stream_name="test-stream",
+            aws_conn_id="aws_test",
+        )
+        classpath, kwargs = trigger.serialize()
+        assert classpath == 
"airflow.providers.amazon.aws.triggers.kinesis.KinesisTrigger"
+        assert kwargs["stream_name"] == "test-stream"
+        assert kwargs["aws_conn_id"] == "aws_test"
+
+    @pytest.mark.asyncio
+    @pytest.mark.usefixtures("cleanup_providers_manager")
+    async def test_message_queue_trigger_run_yields_events(self):
+        from airflow.providers.amazon.aws.triggers.kinesis import 
KinesisTrigger
+        from airflow.providers.common.messaging.triggers.msg_queue import 
MessageQueueTrigger
+        from airflow.triggers.base import TriggerEvent
+
+        trigger = MessageQueueTrigger(
+            scheme="kinesis",
+            stream_name="test-stream",
+            aws_conn_id="aws_test",
+        )
+
+        sample_event = TriggerEvent(
+            {
+                "status": "success",
+                "message_batch": [
+                    {
+                        "ShardId": "shardId-000000000000",
+                        "SequenceNumber": "1",
+                        "PartitionKey": "partition_key_1",
+                        "ApproximateArrivalTimestamp": None,
+                        "Data": "dGVzdF9kYXRh",
+                    }
+                ],
+            }
+        )
+
+        async def mock_run():
+            yield sample_event
+
+        with mock.patch.object(KinesisTrigger, "run", return_value=mock_run()):
+            events = []
+            async for event in trigger.run():
+                events.append(event)
+
+            assert len(events) == 1
+            assert events[0].payload == sample_event.payload

Reply via email to