diff --git a/providers/amazon/docs/message-queues/index.rst b/providers/amazon/docs/message-queues/index.rst index 2d07a423afa66..3d5af93769db9 100644 --- a/providers/amazon/docs/message-queues/index.rst +++ b/providers/amazon/docs/message-queues/index.rst @@ -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` diff --git a/providers/amazon/src/airflow/providers/amazon/aws/queues/sqs.py b/providers/amazon/src/airflow/providers/amazon/aws/queues/sqs.py index 6d2222620716d..910807aefeace 100644 --- a/providers/amazon/src/airflow/providers/amazon/aws/queues/sqs.py +++ b/providers/amazon/src/airflow/providers/amazon/aws/queues/sqs.py @@ -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" diff --git a/providers/apache/kafka/docs/message-queues/index.rst b/providers/apache/kafka/docs/message-queues/index.rst index 13482f4c484ad..8bf6d271b59ef 100644 --- a/providers/apache/kafka/docs/message-queues/index.rst +++ b/providers/apache/kafka/docs/message-queues/index.rst @@ -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: diff --git a/providers/apache/kafka/src/airflow/providers/apache/kafka/queues/kafka.py b/providers/apache/kafka/src/airflow/providers/apache/kafka/queues/kafka.py index c5f7cebaec58e..bb1b1fe3eacc4 100644 --- a/providers/apache/kafka/src/airflow/providers/apache/kafka/queues/kafka.py +++ b/providers/apache/kafka/src/airflow/providers/apache/kafka/queues/kafka.py @@ -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" diff --git a/providers/google/docs/message-queues/index.rst b/providers/google/docs/message-queues/index.rst index fe763999a1807..cae25086c60a0 100644 --- a/providers/google/docs/message-queues/index.rst +++ b/providers/google/docs/message-queues/index.rst @@ -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 ----------------------------- diff --git a/providers/google/src/airflow/providers/google/event_scheduling/events/pubsub.py b/providers/google/src/airflow/providers/google/event_scheduling/events/pubsub.py index 3c808e9fc6b56..ca11c2bf8cbd9 100644 --- a/providers/google/src/airflow/providers/google/event_scheduling/events/pubsub.py +++ b/providers/google/src/airflow/providers/google/event_scheduling/events/pubsub.py @@ -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" diff --git a/providers/ibm/mq/docs/message-queues.rst b/providers/ibm/mq/docs/message-queues.rst index 076c8d92bb19a..3893287ad2a52 100644 --- a/providers/ibm/mq/docs/message-queues.rst +++ b/providers/ibm/mq/docs/message-queues.rst @@ -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: diff --git a/providers/ibm/mq/src/airflow/providers/ibm/mq/queues/mq.py b/providers/ibm/mq/src/airflow/providers/ibm/mq/queues/mq.py index 3bb6382d495e6..342143c872248 100644 --- a/providers/ibm/mq/src/airflow/providers/ibm/mq/queues/mq.py +++ b/providers/ibm/mq/src/airflow/providers/ibm/mq/queues/mq.py @@ -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" diff --git a/providers/microsoft/azure/docs/message-queues/index.rst b/providers/microsoft/azure/docs/message-queues/index.rst index a0a9fba9ac2b5..bbcfef71fea38 100644 --- a/providers/microsoft/azure/docs/message-queues/index.rst +++ b/providers/microsoft/azure/docs/message-queues/index.rst @@ -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 ----------------------------------------- diff --git a/providers/microsoft/azure/src/airflow/providers/microsoft/azure/queues/asb.py b/providers/microsoft/azure/src/airflow/providers/microsoft/azure/queues/asb.py index 867be6693f7d1..065e2468ca1b3 100644 --- a/providers/microsoft/azure/src/airflow/providers/microsoft/azure/queues/asb.py +++ b/providers/microsoft/azure/src/airflow/providers/microsoft/azure/queues/asb.py @@ -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" diff --git a/providers/redis/docs/message-queues.rst b/providers/redis/docs/message-queues.rst index f35aba1a9ff3c..58ee2390cdf9d 100644 --- a/providers/redis/docs/message-queues.rst +++ b/providers/redis/docs/message-queues.rst @@ -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: diff --git a/providers/redis/src/airflow/providers/redis/queues/redis.py b/providers/redis/src/airflow/providers/redis/queues/redis.py index c6ec22b98bb4a..c4693ac73e006 100644 --- a/providers/redis/src/airflow/providers/redis/queues/redis.py +++ b/providers/redis/src/airflow/providers/redis/queues/redis.py @@ -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"