Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 19 additions & 3 deletions providers/amazon/docs/message-queues/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,22 @@ 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 :class:`~airflow.providers.common.messaging.triggers.msg_queue.MessageQueueTrigger` common message queue interface.


.. include:: /../src/airflow/providers/amazon/aws/queues/sqs.py
:start-after: [START sqs_message_queue_provider_description]
:end-before: [END sqs_message_queue_provider_description]
* It uses ``sqs`` as scheme for identifying SQS queues.
* For parameter definitions take a look at :class:`~airflow.providers.amazon.aws.triggers.sqs.SqsSensorTrigger`.

.. code-block:: python

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:
:mod:`tests.system.amazon.aws.example_dag_sqs_message_queue_trigger`
26 changes: 3 additions & 23 deletions providers/amazon/src/airflow/providers/amazon/aws/queues/sqs.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,29 +41,9 @@ class SqsMessageQueueProvider(BaseMessageQueueProvider):
"""
Configuration for SQS integration with common-messaging.

[START sqs_message_queue_provider_description]

* It uses ``sqs`` as scheme for identifying SQS queues.
* For parameter definitions take a look at :class:`~airflow.providers.amazon.aws.triggers.sqs.SqsSensorTrigger`.

.. code-block:: python

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:
:mod:`tests.system.amazon.aws.example_dag_sqs_message_queue_trigger`

[END sqs_message_queue_provider_description]
Dispatches ``scheme="sqs"`` to
:class:`~airflow.providers.amazon.aws.triggers.sqs.SqsSensorTrigger`,
which also defines the accepted parameters.
"""

scheme = "sqs"
Expand Down
23 changes: 20 additions & 3 deletions providers/apache/kafka/docs/message-queues/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -35,9 +35,26 @@ Apache Kafka as the underlying message queue system.
It allows you to send and receive messages using Kafka topics in your Airflow workflows with :class:`~airflow.providers.common.messaging.triggers.msg_queue.MessageQueueTrigger` common message queue interface.


.. include:: /../src/airflow/providers/apache/kafka/queues/kafka.py
:start-after: [START kafka_message_queue_provider_description]
:end-before: [END kafka_message_queue_provider_description]
* It uses ``kafka`` as scheme for identifying Kafka queues.
* For parameter definitions take a look at :class:`~airflow.providers.apache.kafka.triggers.await_message.AwaitMessageTrigger`.

.. code-block:: python

from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger
from airflow.sdk import Asset, AssetWatcher

trigger = MessageQueueTrigger(
scheme="kafka",
# Additional Kafka AwaitMessageTrigger parameters as needed
topics=["my_topic"],
apply_function="module.apply_function",
bootstrap_servers="localhost:9092",
)

asset = Asset("kafka_queue_asset", watchers=[AssetWatcher(name="kafka_watcher", trigger=trigger)])

For a complete example, see:
:mod:`tests.system.common.messaging.kafka_message_queue_trigger`


.. _howto/triggers:KafkaMessageQueueTrigger:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,30 +35,9 @@ class KafkaMessageQueueProvider(BaseMessageQueueProvider):
"""
Configuration for Apache Kafka integration with common-messaging.

[START kafka_message_queue_provider_description]

* It uses ``kafka`` as scheme for identifying Kafka queues.
* For parameter definitions take a look at :class:`~airflow.providers.apache.kafka.triggers.await_message.AwaitMessageTrigger`.

.. code-block:: python

from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger
from airflow.sdk import Asset, AssetWatcher

trigger = MessageQueueTrigger(
scheme="kafka",
# Additional Kafka AwaitMessageTrigger parameters as needed
topics=["my_topic"],
apply_function="module.apply_function",
bootstrap_servers="localhost:9092",
)

asset = Asset("kafka_queue_asset", watchers=[AssetWatcher(name="kafka_watcher", trigger=trigger)])

For a complete example, see:
:mod:`tests.system.common.messaging.kafka_message_queue_trigger`

[END kafka_message_queue_provider_description]
Dispatches ``scheme="kafka"`` to
:class:`~airflow.providers.apache.kafka.triggers.await_message.AwaitMessageTrigger`,
which also defines the accepted parameters.
"""

scheme = "kafka"
Expand Down
24 changes: 21 additions & 3 deletions providers/google/docs/message-queues/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -32,9 +32,27 @@ The Google Cloud Pub/Sub Queue Provider is a message queue provider that uses Go
It allows you to send and receive messages using Cloud Pub/Sub in your Airflow workflows
with :class:`~airflow.providers.common.messaging.triggers.msg_queue.MessageQueueTrigger` common message queue interface.

.. include:: /../src/airflow/providers/google/event_scheduling/events/pubsub.py
:start-after: [START pubsub_message_queue_provider_description]
:end-before: [END pubsub_message_queue_provider_description]
* It uses ``google+pubsub`` as the scheme for identifying the provider.
* For parameter definitions, take a look at :class:`~airflow.providers.google.cloud.triggers.pubsub.PubsubPullTrigger`.

.. code-block:: python

from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger
from airflow.sdk import Asset, AssetWatcher

trigger = MessageQueueTrigger(
scheme="google+pubsub",
# Additional PubsubPullTrigger parameters as needed
project_id="my_project",
subscription="my_subscription",
ack_messages=True,
max_messages=1,
gcp_conn_id="google_cloud_default",
poke_interval=60.0,
return_immediately=False,
)

asset = Asset("pubsub_queue_asset", watchers=[AssetWatcher(name="pubsub_watcher", trigger=trigger)])

Pub/Sub Message Queue Trigger
-----------------------------
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,31 +29,9 @@ class PubSubMessageQueueEventTriggerContainer(BaseMessageQueueProvider):
"""
Configuration for PubSub integration with common-messaging.

[START pubsub_message_queue_provider_description]
* It uses ``google+pubsub`` as the scheme for identifying the provider.
* For parameter definitions, take a look at :class:`~airflow.providers.google.cloud.triggers.pubsub.PubsubPullTrigger`.

.. code-block:: python

from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger
from airflow.sdk import Asset, AssetWatcher

trigger = MessageQueueTrigger(
scheme="google+pubsub",
# Additional PubsubPullTrigger parameters as needed
project_id="my_project",
subscription="my_subscription",
ack_messages=True,
max_messages=1,
gcp_conn_id="google_cloud_default",
poke_interval=60.0,
return_immediately=False,
)

asset = Asset("pubsub_queue_asset", watchers=[AssetWatcher(name="pubsub_watcher", trigger=trigger)])

[END pubsub_message_queue_provider_description]

Dispatches ``scheme="google+pubsub"`` to
:class:`~airflow.providers.google.cloud.triggers.pubsub.PubsubPullTrigger`,
which also defines the accepted parameters.
"""

scheme = "google+pubsub"
Expand Down
17 changes: 14 additions & 3 deletions providers/ibm/mq/docs/message-queues.rst
Original file line number Diff line number Diff line change
Expand Up @@ -43,9 +43,20 @@ via the common message queue interface
:class:`~airflow.providers.common.messaging.triggers.msg_queue.MessageQueueTrigger`.


.. include:: /../src/airflow/providers/ibm/mq/queues/mq.py
:start-after: [START ibmmq_message_queue_provider_description]
:end-before: [END ibmmq_message_queue_provider_description]
* It uses ``ibmmq`` as scheme for identifying IBM MQ queues.
* For parameter definitions take a look at
:class:`~airflow.providers.ibm.mq.triggers.mq.AwaitMessageTrigger`.

.. code-block:: python

from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger
from airflow.sdk import Asset, AssetWatcher

trigger = MessageQueueTrigger(
queue="ibmmq://mq_default/MY.QUEUE.NAME",
)

asset = Asset("mq_topic_asset", watchers=[AssetWatcher(name="mq_watcher", trigger=trigger)])


.. _howto/triggers:IBMMQMessageQueueTrigger:
Expand Down
21 changes: 3 additions & 18 deletions providers/ibm/mq/src/airflow/providers/ibm/mq/queues/mq.py
Original file line number Diff line number Diff line change
Expand Up @@ -37,24 +37,9 @@ class IBMMQMessageQueueProvider(BaseMessageQueueProvider):
"""
Configuration for IBM MQ integration with common-messaging.

[START ibmmq_message_queue_provider_description]

* It uses ``ibmmq`` as scheme for identifying IBM MQ queues.
* For parameter definitions take a look at
:class:`~airflow.providers.ibm.mq.triggers.mq.AwaitMessageTrigger`.

.. code-block:: python

from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger
from airflow.sdk import Asset, AssetWatcher

trigger = MessageQueueTrigger(
queue="ibmmq://mq_default/MY.QUEUE.NAME",
)

asset = Asset("mq_topic_asset", watchers=[AssetWatcher(name="mq_watcher", trigger=trigger)])

[END ibmmq_message_queue_provider_description]
Dispatches ``scheme="ibmmq"`` to
:class:`~airflow.providers.ibm.mq.triggers.mq.AwaitMessageTrigger`,
which also defines the accepted parameters.
"""

scheme = "ibmmq"
Expand Down
24 changes: 21 additions & 3 deletions providers/microsoft/azure/docs/message-queues/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -32,9 +32,27 @@ Azure Service Bus as the underlying message queue system.
It allows you to send and receive messages using Azure Service Bus queues in your Airflow workflows with :class:`~airflow.providers.common.messaging.triggers.msg_queue.MessageQueueTrigger` common message queue interface.


.. include:: /../src/airflow/providers/microsoft/azure/queues/asb.py
:start-after: [START azure_servicebus_message_queue_provider_description]
:end-before: [END azure_servicebus_message_queue_provider_description]
* It uses ``azure+servicebus`` as the scheme for identifying the provider.
* For parameter definitions, take a look at
:class:`~airflow.providers.microsoft.azure.triggers.message_bus.AzureServiceBusQueueTrigger`.

.. code-block:: python

from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger
from airflow.sdk import Asset, AssetWatcher

trigger = MessageQueueTrigger(
scheme="azure+servicebus",
# AzureServiceBusQueueTrigger parameters
queues=["my-queue"],
azure_service_bus_conn_id="azure_service_bus_default",
poll_interval=60,
)

asset = Asset(
"asb_queue_asset",
watchers=[AssetWatcher(name="asb_watcher", trigger=trigger)],
)

Azure Service Bus Message Queue Trigger
-----------------------------------------
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,31 +29,9 @@ class AzureServiceBusMessageQueueProvider(BaseMessageQueueProvider):
"""
Configuration for Azure Service Bus integration with common-messaging.

[START azure_servicebus_message_queue_provider_description]

* It uses ``azure+servicebus`` as the scheme for identifying the provider.
* For parameter definitions, take a look at
:class:`~airflow.providers.microsoft.azure.triggers.message_bus.AzureServiceBusQueueTrigger`.

.. code-block:: python

from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger
from airflow.sdk import Asset, AssetWatcher

trigger = MessageQueueTrigger(
scheme="azure+servicebus",
# AzureServiceBusQueueTrigger parameters
queues=["my-queue"],
azure_service_bus_conn_id="azure_service_bus_default",
poll_interval=60,
)

asset = Asset(
"asb_queue_asset",
watchers=[AssetWatcher(name="asb_watcher", trigger=trigger)],
)

[END azure_servicebus_message_queue_provider_description]
Dispatches ``scheme="azure+servicebus"`` to
:class:`~airflow.providers.microsoft.azure.triggers.message_bus.AzureServiceBusQueueTrigger`,
which also defines the accepted parameters.
"""

scheme = "azure+servicebus"
Expand Down
19 changes: 16 additions & 3 deletions providers/redis/docs/message-queues.rst
Original file line number Diff line number Diff line change
Expand Up @@ -39,9 +39,22 @@ Redis as the underlying message queue system.
It allows you to send and receive messages using Redis channels in your Airflow workflows with :class:`~airflow.providers.common.messaging.triggers.msg_queue.MessageQueueTrigger` common message queue interface.


.. include:: /../src/airflow/providers/redis/queues/redis.py
:start-after: [START redis_message_queue_provider_description]
:end-before: [END redis_message_queue_provider_description]
* It uses ``redis+pubsub`` as scheme for identifying Redis queues.
* For parameter definitions take a look at :class:`~airflow.providers.redis.triggers.redis_await_message.AwaitMessageTrigger`.

.. code-block:: python

from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger
from airflow.sdk import Asset, AssetWatcher

trigger = MessageQueueTrigger(
scheme="redis+pubsub",
# Additional Redis AwaitMessageTrigger parameters as needed
channels=["my_channel"],
redis_conn_id="redis_default",
)

asset = Asset("redis_queue_asset", watchers=[AssetWatcher(name="redis_watcher", trigger=trigger)])


.. _howto/triggers:RedisMessageQueueTrigger:
Expand Down
23 changes: 3 additions & 20 deletions providers/redis/src/airflow/providers/redis/queues/redis.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,26 +33,9 @@ class RedisPubSubMessageQueueProvider(BaseMessageQueueProvider):
"""
Configuration for Redis integration with common-messaging.

[START redis_message_queue_provider_description]

* It uses ``redis+pubsub`` as scheme for identifying Redis queues.
* For parameter definitions take a look at :class:`~airflow.providers.redis.triggers.redis_await_message.AwaitMessageTrigger`.

.. code-block:: python

from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger
from airflow.sdk import Asset, AssetWatcher

trigger = MessageQueueTrigger(
scheme="redis+pubsub",
# Additional Redis AwaitMessageTrigger parameters as needed
channels=["my_channel"],
redis_conn_id="redis_default",
)

asset = Asset("redis_queue_asset", watchers=[AssetWatcher(name="redis_watcher", trigger=trigger)])

[END redis_message_queue_provider_description]
Dispatches ``scheme="redis+pubsub"`` to
:class:`~airflow.providers.redis.triggers.redis_await_message.AwaitMessageTrigger`,
which also defines the accepted parameters.
"""

scheme = "redis+pubsub"
Expand Down