This is an automated email from the ASF dual-hosted git repository.
potiuk 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 ea03b2921c3 Make Pub/Sub return_immediately configurable, deprecate
its default (#73504)
ea03b2921c3 is described below
commit ea03b2921c3c65d3900e9f9175ef54f48c567468
Author: Shahar Epstein <[email protected]>
AuthorDate: Tue Sep 22 12:42:06 2026 +0300
Make Pub/Sub return_immediately configurable, deprecate its default (#73504)
* Make return_immediately configurable for the Pub/Sub modules
The Pub/Sub Pull operator, sensor and trigger hardcoded
return_immediately=True, which relies on a Pull option Google deprecated
because it can return zero messages while a backlog exists. Users had no way
to opt into the long-polling behaviour Google recommends instead.
Keeping True as the default preserves existing behaviour, so the change is
paired with a deprecation warning announcing the coming flip.
* Warn when return_immediately is unset rather than when it is set
Warning only when a user passes the option leaves the people who most need
the notice -- everyone still on the implicit default -- hearing nothing, and
it nags the users who already made a deliberate choice. Inverting it also
lets the argument stay absent from a call without silently taking the
deprecated path.
PubsubPullTrigger gains the same treatment. Nothing covered it before, yet
it
is constructed directly by the google+pubsub scheme for asset watchers, so
those Dag authors were getting the deprecated default with nothing telling
them. It names the subscription in its message because it is built inside
MessageQueueTrigger.serialize(), where stacklevel=2 resolves to
common.messaging's file rather than the user's watcher. The operator and the
sensor cannot do the same: subscription is a template field there, so at
__init__ time it can still hold an unrendered Jinja expression.
Only the message text is shared, in one constant. The warn call stays in
each
class: fixup_decorator_warning_stack only adjusts the stack for modules that
define an operator, so moving the call out of them would break the frame the
warning points at.
The tests drop a module-level filterwarnings mark that silenced every
return_immediately deprecation in these files, and start asserting that an
unset argument still resolves to True -- the backward-compatibility contract
of the deprecation, which nothing pinned.
* Set return_immediately in the Pub/Sub examples
MessageQueueTrigger.serialize() builds a PubsubPullTrigger from the keyword
arguments it was given, so a google+pubsub watcher that omits
return_immediately warns during Dag serialization, in the Dag processor
rather than in any task log. The examples are rendered into the operator and
message-queues guides, so they were teaching the pattern the deprecation
steers users away from.
* Document the Pub/Sub return_immediately changes for users
Two things reach users and neither is visible from the release notes
otherwise: the new deprecation warning, which for asset watchers appears in
Dag processor logs where nobody looks for it, and the deferrable sensor
starting to honour an argument it used to drop on the floor. The operator
guide also never said that the operator waits indefinitely in deferrable
mode, which is the behaviour most likely to surprise.
---------
Co-authored-by: Mike Prieto <[email protected]>
---
providers/google/docs/changelog.rst | 14 +++++
providers/google/docs/operators/cloud/pubsub.rst | 7 +++
.../providers/google/cloud/operators/pubsub.py | 38 +++++++++++--
.../providers/google/cloud/sensors/pubsub.py | 33 +++++++----
.../providers/google/cloud/triggers/pubsub.py | 24 +++++++-
.../src/airflow/providers/google/common/consts.py | 8 +++
.../google/event_scheduling/events/pubsub.py | 1 +
.../system/google/cloud/pubsub/example_pubsub.py | 5 ++
.../cloud/pubsub/example_pubsub_deferrable.py | 1 +
.../example_event_schedule_pubsub.py | 1 +
.../unit/google/cloud/operators/test_pubsub.py | 47 ++++++++++++++-
.../tests/unit/google/cloud/sensors/test_pubsub.py | 45 +++++++++++++--
.../unit/google/cloud/triggers/test_pubsub.py | 66 ++++++++++++++++++++++
13 files changed, 266 insertions(+), 24 deletions(-)
diff --git a/providers/google/docs/changelog.rst
b/providers/google/docs/changelog.rst
index 5bb751bb5e3..86ac6b0cf7e 100644
--- a/providers/google/docs/changelog.rst
+++ b/providers/google/docs/changelog.rst
@@ -27,6 +27,20 @@
Changelog
---------
+.. note::
+ ``PubSubPullOperator``, ``PubSubPullSensor`` and ``PubsubPullTrigger`` now
emit a deprecation
+ warning when ``return_immediately`` is left unset -- including
``google+pubsub`` asset
+ watchers built with ``MessageQueueTrigger``, where the warning surfaces in
Dag processor
+ logs rather than task logs. It currently defaults to ``True``, which
relies on the deprecated
+ Pub/Sub ``returnImmediately`` Pull option and can return zero messages
even when a backlog
+ exists. The default will change to ``False`` in the first Google provider
major release after
+ March 31, 2027 -- pass ``return_immediately=True`` explicitly to keep the
current behaviour.
+
+ Deferrable ``PubSubPullSensor`` now respects ``return_immediately`` as
well. It previously
+ dropped the argument when handing off to ``PubsubPullTrigger``, so the
trigger always behaved
+ as ``True``. A Dag already using ``PubSubPullSensor(deferrable=True,
return_immediately=False)``
+ will see its triggerer start long-polling on each pull instead of
returning immediately.
+
22.5.0
......
diff --git a/providers/google/docs/operators/cloud/pubsub.rst
b/providers/google/docs/operators/cloud/pubsub.rst
index a113b2e775f..2376dcbe291 100644
--- a/providers/google/docs/operators/cloud/pubsub.rst
+++ b/providers/google/docs/operators/cloud/pubsub.rst
@@ -96,6 +96,13 @@ Also for this action you can use sensor in the deferrable
mode:
:start-after: [START
howto_operator_gcp_pubsub_pull_message_with_async_sensor]
:end-before: [END howto_operator_gcp_pubsub_pull_message_with_async_sensor]
+Unlike the sensor, which pokes until a message shows up, the
+:class:`~airflow.providers.google.cloud.operators.pubsub.PubSubPullOperator`
operator does not poke.
+With ``return_immediately=True`` it issues a single pull, and an empty
subscription yields an empty
+list. In deferrable mode it hands the wait to
+:class:`~airflow.providers.google.cloud.triggers.pubsub.PubsubPullTrigger`,
which re-pulls every
+``poll_interval`` until a message arrives, with nothing bounding that wait.
+
.. exampleinclude::
/../../google/tests/system/google/cloud/pubsub/example_pubsub.py
:language: python
:start-after: [START howto_operator_gcp_pubsub_pull_message_with_operator]
diff --git
a/providers/google/src/airflow/providers/google/cloud/operators/pubsub.py
b/providers/google/src/airflow/providers/google/cloud/operators/pubsub.py
index 33a1b302786..6caad672cda 100644
--- a/providers/google/src/airflow/providers/google/cloud/operators/pubsub.py
+++ b/providers/google/src/airflow/providers/google/cloud/operators/pubsub.py
@@ -25,6 +25,7 @@ This module contains Google PubSub operators.
from __future__ import annotations
+import warnings
from collections.abc import Callable, Sequence
from functools import cached_property
from typing import TYPE_CHECKING, Any
@@ -42,12 +43,16 @@ from google.cloud.pubsub_v1.types import (
SchemaSettings,
)
+from airflow.exceptions import AirflowProviderDeprecationWarning
from airflow.providers.common.compat.sdk import AirflowException, conf
from airflow.providers.google.cloud.hooks.pubsub import PubSubHook
from airflow.providers.google.cloud.links.pubsub import
PubSubSubscriptionLink, PubSubTopicLink
from airflow.providers.google.cloud.operators.cloud_base import
GoogleCloudBaseOperator
from airflow.providers.google.cloud.triggers.pubsub import PubsubPullTrigger
-from airflow.providers.google.common.consts import
GOOGLE_DEFAULT_DEFERRABLE_METHOD_NAME
+from airflow.providers.google.common.consts import (
+ GOOGLE_DEFAULT_DEFERRABLE_METHOD_NAME,
+ PUBSUB_RETURN_IMMEDIATELY_DEPRECATION_MESSAGE,
+)
from airflow.providers.google.common.hooks.base_google import
PROVIDE_PROJECT_ID
if TYPE_CHECKING:
@@ -756,9 +761,16 @@ class PubSubPullOperator(GoogleCloudBaseOperator):
"""
Pulls messages from a PubSub subscription and passes them through XCom.
- If the queue is empty, returns empty list - never waits for messages.
- If you do need to wait, please use
:class:`airflow.providers.google.cloud.sensors.PubSubPullSensor`
- instead.
+ In non-deferrable mode, ``return_immediately=True`` returns an empty list
when the
+ queue is empty; ``return_immediately=False`` makes the Pub/Sub API block
for a bounded,
+ server-side period for at least one message instead, occupying the worker
slot for that
+ duration. In deferrable mode the operator always waits for at least one
message no matter how
+ ``return_immediately`` is set:
+ :class:`~airflow.providers.google.cloud.triggers.pubsub.PubsubPullTrigger`
re-pulls every
+ ``poll_interval`` until messages arrive — nothing in the operator bounds
that wait — and
+ ``return_immediately`` only controls whether each individual pull
long-polls. For the
+ poke-based equivalent of this waiting behavior, see
+ :class:`~airflow.providers.google.cloud.sensors.pubsub.PubSubPullSensor`.
.. seealso::
For more information on how to use this operator and the
PubSubPullSensor, take a look at the guide:
@@ -799,6 +811,12 @@ class PubSubPullOperator(GoogleCloudBaseOperator):
:param deferrable: If True, run the task in the deferrable mode.
:param poll_interval: Time (seconds) to wait between two consecutive calls
to check the job.
The default is 300 seconds.
+ :param return_immediately: Defaults to True, which uses the deprecated
Pub/Sub
+ ``returnImmediately`` Pull option and can return zero messages even if
there are
+ messages in the backlog. If set to False, the system will instead wait
(for a bounded
+ amount of time) until at least one message is available, rather than
returning no
+ messages. The default will change to False in the first Google
provider major release
+ after March 31, 2027.
"""
template_fields: Sequence[str] = (
@@ -819,6 +837,7 @@ class PubSubPullOperator(GoogleCloudBaseOperator):
impersonation_chain: str | Sequence[str] | None = None,
deferrable: bool = conf.getboolean("operators", "default_deferrable",
fallback=False),
poll_interval: int = 300,
+ return_immediately: bool | None = None,
**kwargs,
) -> None:
super().__init__(**kwargs)
@@ -831,6 +850,14 @@ class PubSubPullOperator(GoogleCloudBaseOperator):
self.impersonation_chain = impersonation_chain
self.deferrable = deferrable
self.poll_interval = poll_interval
+ if return_immediately is None:
+ warnings.warn(
+ PUBSUB_RETURN_IMMEDIATELY_DEPRECATION_MESSAGE,
+ AirflowProviderDeprecationWarning,
+ stacklevel=2,
+ )
+ return_immediately = True
+ self.return_immediately = return_immediately
def execute(self, context: Context) -> list:
if self.deferrable:
@@ -843,6 +870,7 @@ class PubSubPullOperator(GoogleCloudBaseOperator):
gcp_conn_id=self.gcp_conn_id,
poke_interval=self.poll_interval,
impersonation_chain=self.impersonation_chain,
+ return_immediately=self.return_immediately,
),
method_name=GOOGLE_DEFAULT_DEFERRABLE_METHOD_NAME,
)
@@ -855,7 +883,7 @@ class PubSubPullOperator(GoogleCloudBaseOperator):
project_id=self.project_id,
subscription=self.subscription,
max_messages=self.max_messages,
- return_immediately=True,
+ return_immediately=self.return_immediately,
)
handle_messages = self.messages_callback or
self._default_message_callback
diff --git
a/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py
b/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py
index f138271b66e..5e993d5be46 100644
--- a/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py
+++ b/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py
@@ -19,6 +19,7 @@
from __future__ import annotations
+import warnings
from collections.abc import Callable, Sequence
from datetime import timedelta
from typing import TYPE_CHECKING, Any
@@ -26,9 +27,11 @@ from typing import TYPE_CHECKING, Any
from google.cloud import pubsub_v1
from google.cloud.pubsub_v1.types import ReceivedMessage
+from airflow.exceptions import AirflowProviderDeprecationWarning
from airflow.providers.common.compat.sdk import AirflowException,
BaseSensorOperator, conf
from airflow.providers.google.cloud.hooks.pubsub import PubSubHook
from airflow.providers.google.cloud.triggers.pubsub import PubsubPullTrigger
+from airflow.providers.google.common.consts import
PUBSUB_RETURN_IMMEDIATELY_DEPRECATION_MESSAGE
if TYPE_CHECKING:
from airflow.providers.common.compat.sdk import Context
@@ -49,7 +52,8 @@ class PubSubPullSensor(BaseSensorOperator):
:ref:`howto/operator:PubSubPullSensor`
.. seealso::
- If you don't want to wait for at least one message to come, use
Operator instead:
+ If you don't want to wait for at least one message to come, use the
operator with
+ ``return_immediately=True`` and ``deferrable=False`` instead:
:class:`~airflow.providers.google.cloud.operators.pubsub.PubSubPullOperator`
This sensor operator will pull up to ``max_messages`` messages from the
@@ -61,9 +65,9 @@ class PubSubPullSensor(BaseSensorOperator):
acknowledged before being returned, otherwise, downstream tasks will be
responsible for acknowledging them.
- If you want a non-blocking task that does not to wait for messages, please
use
+ If you want a non-blocking task that does not wait for messages, please use
:class:`~airflow.providers.google.cloud.operators.pubsub.PubSubPullOperator`
- instead.
+ with ``return_immediately=True`` and ``deferrable=False`` instead.
``project_id`` and ``subscription`` are templated so you can use
variables in them.
@@ -73,13 +77,12 @@ class PubSubPullSensor(BaseSensorOperator):
full subscription path.
:param max_messages: The maximum number of messages to retrieve per
PubSub pull request
- :param return_immediately: If this field set to true, the system will
- respond immediately even if it there are no messages available to
- return in the ``Pull`` response. Otherwise, the system may wait
- (for a bounded amount of time) until at least one message is available,
- rather than returning no messages. Warning: setting this field to
- ``true`` is discouraged because it adversely impacts the performance
- of ``Pull`` operations. We recommend that users do not set this field.
+ :param return_immediately: Defaults to True, which uses the deprecated
Pub/Sub
+ ``returnImmediately`` Pull option and can return zero messages even if
there are
+ messages in the backlog. If set to False, the system will instead wait
(for a bounded
+ amount of time) until at least one message is available, rather than
returning no
+ messages. The default will change to False in the first Google
provider major release
+ after March 31, 2027.
:param ack_messages: If True, each message will be acknowledged
immediately rather than by any downstream tasks
:param gcp_conn_id: The connection ID to use connecting to
@@ -113,7 +116,7 @@ class PubSubPullSensor(BaseSensorOperator):
project_id: str,
subscription: str,
max_messages: int = 5,
- return_immediately: bool = True,
+ return_immediately: bool | None = None,
ack_messages: bool = False,
gcp_conn_id: str = "google_cloud_default",
messages_callback: Callable[[list[ReceivedMessage], Context], Any] |
None = None,
@@ -127,6 +130,13 @@ class PubSubPullSensor(BaseSensorOperator):
self.project_id = project_id
self.subscription = subscription
self.max_messages = max_messages
+ if return_immediately is None:
+ warnings.warn(
+ PUBSUB_RETURN_IMMEDIATELY_DEPRECATION_MESSAGE,
+ AirflowProviderDeprecationWarning,
+ stacklevel=2,
+ )
+ return_immediately = True
self.return_immediately = return_immediately
self.ack_messages = ack_messages
self.messages_callback = messages_callback
@@ -176,6 +186,7 @@ class PubSubPullSensor(BaseSensorOperator):
poke_interval=self.poke_interval,
gcp_conn_id=self.gcp_conn_id,
impersonation_chain=self.impersonation_chain,
+ return_immediately=self.return_immediately,
),
method_name="execute_complete",
)
diff --git
a/providers/google/src/airflow/providers/google/cloud/triggers/pubsub.py
b/providers/google/src/airflow/providers/google/cloud/triggers/pubsub.py
index 22314838c46..eabe74055cf 100644
--- a/providers/google/src/airflow/providers/google/cloud/triggers/pubsub.py
+++ b/providers/google/src/airflow/providers/google/cloud/triggers/pubsub.py
@@ -19,13 +19,16 @@
from __future__ import annotations
import asyncio
+import warnings
from collections.abc import AsyncIterator, Sequence
from functools import cached_property
from typing import Any
from google.cloud.pubsub_v1.types import ReceivedMessage
+from airflow.exceptions import AirflowProviderDeprecationWarning
from airflow.providers.google.cloud.hooks.pubsub import PubSubAsyncHook
+from airflow.providers.google.common.consts import
PUBSUB_RETURN_IMMEDIATELY_DEPRECATION_MESSAGE
from airflow.providers.google.version_compat import AIRFLOW_V_3_0_PLUS
from airflow.triggers.base import TriggerEvent
@@ -55,6 +58,15 @@ class PubsubPullTrigger(BaseEventTrigger):
If set as a sequence, the identities from the list must grant
Service Account Token Creator IAM role to the directly preceding
identity, with first
account from the list granting this role to the originating account
(templated).
+ :param return_immediately: Normally supplied by the sensor or operator
that defers to this
+ trigger; callers constructing the trigger directly (for example via
+
:class:`~airflow.providers.common.messaging.triggers.msg_queue.MessageQueueTrigger`)
can
+ set it themselves. Defaults to True, which uses the deprecated Pub/Sub
+ ``returnImmediately`` Pull option and can return zero messages even if
there are messages
+ in the backlog. If set to False, the system will instead wait (for a
bounded amount of
+ time) until at least one message is available, rather than returning
no messages. The
+ default will change to False in the first Google provider major
release after
+ March 31, 2027.
"""
def __init__(
@@ -66,6 +78,7 @@ class PubsubPullTrigger(BaseEventTrigger):
gcp_conn_id: str,
poke_interval: float = 10.0,
impersonation_chain: str | Sequence[str] | None = None,
+ return_immediately: bool | None = None,
):
super().__init__()
self.project_id = project_id
@@ -75,6 +88,14 @@ class PubsubPullTrigger(BaseEventTrigger):
self.poke_interval = poke_interval
self.gcp_conn_id = gcp_conn_id
self.impersonation_chain = impersonation_chain
+ if return_immediately is None:
+ warnings.warn(
+ f"{PUBSUB_RETURN_IMMEDIATELY_DEPRECATION_MESSAGE}
Subscription: {self.subscription}.",
+ AirflowProviderDeprecationWarning,
+ stacklevel=2,
+ )
+ return_immediately = True
+ self.return_immediately = return_immediately
def serialize(self) -> tuple[str, dict[str, Any]]:
"""Serialize PubsubPullTrigger arguments and classpath."""
@@ -88,6 +109,7 @@ class PubsubPullTrigger(BaseEventTrigger):
"poke_interval": self.poke_interval,
"gcp_conn_id": self.gcp_conn_id,
"impersonation_chain": self.impersonation_chain,
+ "return_immediately": self.return_immediately,
},
)
@@ -97,7 +119,7 @@ class PubsubPullTrigger(BaseEventTrigger):
project_id=self.project_id,
subscription=self.subscription,
max_messages=self.max_messages,
- return_immediately=True,
+ return_immediately=self.return_immediately,
):
if self.ack_messages:
await self.message_acknowledgement(pulled_messages)
diff --git a/providers/google/src/airflow/providers/google/common/consts.py
b/providers/google/src/airflow/providers/google/common/consts.py
index f8d7209901d..cd7ae3d912b 100644
--- a/providers/google/src/airflow/providers/google/common/consts.py
+++ b/providers/google/src/airflow/providers/google/common/consts.py
@@ -23,3 +23,11 @@ from airflow import version
GOOGLE_DEFAULT_DEFERRABLE_METHOD_NAME = "execute_complete"
CLIENT_INFO = ClientInfo(client_library_version="airflow_v" + version.version)
+
+PUBSUB_RETURN_IMMEDIATELY_DEPRECATION_MESSAGE = (
+ "`return_immediately` defaults to True, which relies on the deprecated
Pub/Sub "
+ "`returnImmediately` Pull option and can return zero messages while a
backlog exists. "
+ "The default will change to False in the first Google provider major
release after "
+ "March 31, 2027. Pass `return_immediately=False` to adopt the new
behaviour now, "
+ "or `return_immediately=True` to keep the current one."
+)
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 1d1e3612181..3c808e9fc6b 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
@@ -47,6 +47,7 @@ class
PubSubMessageQueueEventTriggerContainer(BaseMessageQueueProvider):
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)])
diff --git
a/providers/google/tests/system/google/cloud/pubsub/example_pubsub.py
b/providers/google/tests/system/google/cloud/pubsub/example_pubsub.py
index f2b6dae0158..de6608aad26 100644
--- a/providers/google/tests/system/google/cloud/pubsub/example_pubsub.py
+++ b/providers/google/tests/system/google/cloud/pubsub/example_pubsub.py
@@ -85,6 +85,7 @@ with DAG(
ack_messages=True,
project_id=PROJECT_ID,
subscription=subscription,
+ return_immediately=False,
)
# [END howto_operator_gcp_pubsub_pull_message_with_sensor]
@@ -94,11 +95,15 @@ with DAG(
# [START howto_operator_gcp_pubsub_pull_message_with_operator]
+ # return_immediately=False makes this pull block for a bounded,
server-side period, holding the
+ # worker slot; pass return_immediately=True explicitly if the task should
return an empty list
+ # instead of waiting.
pull_messages_operator = PubSubPullOperator(
task_id="pull_messages_operator",
ack_messages=True,
project_id=PROJECT_ID,
subscription=subscription,
+ return_immediately=False,
)
# [END howto_operator_gcp_pubsub_pull_message_with_operator]
diff --git
a/providers/google/tests/system/google/cloud/pubsub/example_pubsub_deferrable.py
b/providers/google/tests/system/google/cloud/pubsub/example_pubsub_deferrable.py
index 195f5a313a5..78c92f44c5f 100644
---
a/providers/google/tests/system/google/cloud/pubsub/example_pubsub_deferrable.py
+++
b/providers/google/tests/system/google/cloud/pubsub/example_pubsub_deferrable.py
@@ -78,6 +78,7 @@ with DAG(
project_id=PROJECT_ID,
subscription=subscription,
deferrable=True,
+ return_immediately=False,
)
# [END howto_operator_gcp_pubsub_pull_message_with_async_sensor]
diff --git
a/providers/google/tests/system/google/event_scheduling/example_event_schedule_pubsub.py
b/providers/google/tests/system/google/event_scheduling/example_event_schedule_pubsub.py
index 4dbc466cc7d..93a2aa581e7 100644
---
a/providers/google/tests/system/google/event_scheduling/example_event_schedule_pubsub.py
+++
b/providers/google/tests/system/google/event_scheduling/example_event_schedule_pubsub.py
@@ -74,6 +74,7 @@ trigger = MessageQueueTrigger(
max_messages=1,
gcp_conn_id="google_cloud_default",
poke_interval=60.0,
+ return_immediately=False,
)
# Define an asset that watches for messages on the Pub/Sub subscription
diff --git a/providers/google/tests/unit/google/cloud/operators/test_pubsub.py
b/providers/google/tests/unit/google/cloud/operators/test_pubsub.py
index 3537c5266db..946a79d6a81 100644
--- a/providers/google/tests/unit/google/cloud/operators/test_pubsub.py
+++ b/providers/google/tests/unit/google/cloud/operators/test_pubsub.py
@@ -17,6 +17,7 @@
# under the License.
from __future__ import annotations
+import warnings
from typing import Any
from unittest import mock
@@ -25,6 +26,7 @@ from google.api_core.gapic_v1.method import DEFAULT
from google.cloud import pubsub_v1
from google.cloud.pubsub_v1.types import ReceivedMessage
+from airflow.exceptions import AirflowProviderDeprecationWarning
from airflow.providers.common.compat.sdk import TaskDeferred
from airflow.providers.google.cloud.operators.pubsub import (
PubSubCreateSubscriptionOperator,
@@ -34,6 +36,7 @@ from airflow.providers.google.cloud.operators.pubsub import (
PubSubPublishMessageOperator,
PubSubPullOperator,
)
+from airflow.providers.google.cloud.triggers.pubsub import PubsubPullTrigger
TASK_ID = "test-task-id"
TEST_PROJECT = "test-project"
@@ -445,16 +448,24 @@ class TestPubSubPullOperator:
def _generate_dicts(self, count):
return [ReceivedMessage.to_dict(m) for m in
self._generate_messages(count)]
+ @pytest.mark.parametrize("return_immediately", [True, False])
@mock.patch("airflow.providers.google.cloud.operators.pubsub.PubSubHook")
- def test_execute_no_messages(self, mock_hook):
+ def test_execute_no_messages(self, mock_hook, return_immediately):
operator = PubSubPullOperator(
task_id=TASK_ID,
project_id=TEST_PROJECT,
subscription=TEST_SUBSCRIPTION,
+ return_immediately=return_immediately,
)
mock_hook.return_value.pull.return_value = []
assert operator.execute({}) == []
+ mock_hook.return_value.pull.assert_called_once_with(
+ project_id=TEST_PROJECT,
+ subscription=TEST_SUBSCRIPTION,
+ max_messages=5,
+ return_immediately=return_immediately,
+ )
@mock.patch("airflow.providers.google.cloud.operators.pubsub.PubSubHook")
def test_execute_with_ack_messages(self, mock_hook):
@@ -463,6 +474,7 @@ class TestPubSubPullOperator:
project_id=TEST_PROJECT,
subscription=TEST_SUBSCRIPTION,
ack_messages=True,
+ return_immediately=True,
)
generated_messages = self._generate_messages(5)
@@ -500,6 +512,7 @@ class TestPubSubPullOperator:
project_id=TEST_PROJECT,
subscription=TEST_SUBSCRIPTION,
messages_callback=messages_callback,
+ return_immediately=True,
)
mock_hook.return_value.pull.return_value = generated_messages
@@ -514,8 +527,9 @@ class TestPubSubPullOperator:
assert response == messages_callback_return_value
@pytest.mark.db_test
+ @pytest.mark.parametrize("return_immediately", [True, False])
@mock.patch("airflow.providers.google.cloud.operators.pubsub.PubSubHook")
- def test_execute_deferred(self, mock_hook):
+ def test_execute_deferred(self, mock_hook, return_immediately):
"""
Asserts that a task is deferred and a PubSubPullOperator will be fired
when the PubSubPullOperator is executed with deferrable=True.
@@ -525,16 +539,21 @@ class TestPubSubPullOperator:
project_id=TEST_PROJECT,
subscription=TEST_SUBSCRIPTION,
deferrable=True,
+ return_immediately=return_immediately,
)
- with pytest.raises(TaskDeferred) as _:
+ with pytest.raises(TaskDeferred) as exc:
task.execute(mock.MagicMock())
+ assert isinstance(exc.value.trigger, PubsubPullTrigger)
+ assert exc.value.trigger.return_immediately is return_immediately
+
@mock.patch("airflow.providers.google.cloud.operators.pubsub.PubSubHook")
def test_get_openlineage_facets(self, mock_hook):
operator = PubSubPullOperator(
task_id=TASK_ID,
project_id=TEST_PROJECT,
subscription=TEST_SUBSCRIPTION,
+ return_immediately=True,
)
generated_messages = self._generate_messages(5)
@@ -593,6 +612,7 @@ class TestPubSubPullOperator:
subscription=TEST_SUBSCRIPTION,
deferrable=True,
messages_callback=messages_callback,
+ return_immediately=True,
)
mock_hook.return_value.pull.return_value = received_messages
@@ -624,6 +644,7 @@ class TestPubSubPullOperator:
project_id=TEST_PROJECT,
subscription=TEST_SUBSCRIPTION,
deferrable=True,
+ return_immediately=True,
)
mock_hook.return_value.pull.return_value = received_messages
@@ -631,3 +652,23 @@ class TestPubSubPullOperator:
resp = operator.execute_complete(context={}, event={"status":
"success", "message": test_message})
mock_log_info.assert_called_with("Sensor pulls messages: %s",
test_message)
assert resp == [ReceivedMessage.to_dict(m) for m in received_messages]
+
+ def test_pubsub_pull_operator_deprecation_warning(self):
+ with pytest.warns(AirflowProviderDeprecationWarning,
match="return_immediately"):
+ operator = PubSubPullOperator(
+ task_id=TASK_ID,
+ project_id=TEST_PROJECT,
+ subscription=TEST_SUBSCRIPTION,
+ )
+ assert operator.return_immediately is True
+
+ @pytest.mark.parametrize("return_immediately", [True, False])
+ def test_pubsub_pull_operator_no_deprecation_warning_when_explicit(self,
return_immediately):
+ with warnings.catch_warnings():
+ warnings.simplefilter("error", AirflowProviderDeprecationWarning)
+ PubSubPullOperator(
+ task_id=TASK_ID,
+ project_id=TEST_PROJECT,
+ subscription=TEST_SUBSCRIPTION,
+ return_immediately=return_immediately,
+ )
diff --git a/providers/google/tests/unit/google/cloud/sensors/test_pubsub.py
b/providers/google/tests/unit/google/cloud/sensors/test_pubsub.py
index 4cd1b48fbfb..d5f4e692b0f 100644
--- a/providers/google/tests/unit/google/cloud/sensors/test_pubsub.py
+++ b/providers/google/tests/unit/google/cloud/sensors/test_pubsub.py
@@ -17,6 +17,7 @@
# under the License.
from __future__ import annotations
+import warnings
from typing import Any
from unittest import mock
@@ -24,6 +25,7 @@ import pytest
from google.cloud import pubsub_v1
from google.cloud.pubsub_v1.types import ReceivedMessage
+from airflow.exceptions import AirflowProviderDeprecationWarning
from airflow.providers.common.compat.sdk import AirflowException, TaskDeferred
from airflow.providers.google.cloud.sensors.pubsub import PubSubPullSensor
from airflow.providers.google.cloud.triggers.pubsub import PubsubPullTrigger
@@ -55,6 +57,7 @@ class TestPubSubPullSensor:
task_id=TASK_ID,
project_id=TEST_PROJECT,
subscription=TEST_SUBSCRIPTION,
+ return_immediately=True,
)
mock_hook.return_value.pull.return_value = []
@@ -67,6 +70,7 @@ class TestPubSubPullSensor:
project_id=TEST_PROJECT,
subscription=TEST_SUBSCRIPTION,
ack_messages=True,
+ return_immediately=True,
)
generated_messages = self._generate_messages(5)
@@ -80,13 +84,15 @@ class TestPubSubPullSensor:
messages=generated_messages,
)
+ @pytest.mark.parametrize("return_immediately", [True, False])
@mock.patch("airflow.providers.google.cloud.sensors.pubsub.PubSubHook")
- def test_execute(self, mock_hook):
+ def test_execute(self, mock_hook, return_immediately):
operator = PubSubPullSensor(
task_id=TASK_ID,
project_id=TEST_PROJECT,
subscription=TEST_SUBSCRIPTION,
poke_interval=0,
+ return_immediately=return_immediately,
)
generated_messages = self._generate_messages(5)
@@ -95,7 +101,10 @@ class TestPubSubPullSensor:
response = operator.execute({})
mock_hook.return_value.pull.assert_called_once_with(
- project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION,
max_messages=5, return_immediately=True
+ project_id=TEST_PROJECT,
+ subscription=TEST_SUBSCRIPTION,
+ max_messages=5,
+ return_immediately=return_immediately,
)
assert generated_dicts == response
@@ -107,6 +116,7 @@ class TestPubSubPullSensor:
subscription=TEST_SUBSCRIPTION,
poke_interval=0,
timeout=1,
+ return_immediately=True,
)
mock_hook.return_value.pull.return_value = []
@@ -139,6 +149,7 @@ class TestPubSubPullSensor:
subscription=TEST_SUBSCRIPTION,
poke_interval=0,
messages_callback=messages_callback,
+ return_immediately=True,
)
mock_hook.return_value.pull.return_value = generated_messages
@@ -152,10 +163,11 @@ class TestPubSubPullSensor:
assert response == messages_callback_return_value
- def test_pubsub_pull_sensor_async(self):
+ @pytest.mark.parametrize("return_immediately", [True, False])
+ def test_pubsub_pull_sensor_async(self, return_immediately):
"""
Asserts that a task is deferred and a PubsubPullTrigger will be fired
- when the PubSubPullSensor is executed.
+ when the PubSubPullSensor is executed, with the configured
return_immediately value.
"""
task = PubSubPullSensor(
task_id="test_task_id",
@@ -163,10 +175,12 @@ class TestPubSubPullSensor:
project_id=TEST_PROJECT,
subscription=TEST_SUBSCRIPTION,
deferrable=True,
+ return_immediately=return_immediately,
)
with pytest.raises(TaskDeferred) as exc:
task.execute(context={})
assert isinstance(exc.value.trigger, PubsubPullTrigger), "Trigger is
not a PubsubPullTrigger"
+ assert exc.value.trigger.return_immediately is return_immediately
def test_pubsub_pull_sensor_async_execute_should_throw_exception(self):
"""Tests that an AirflowException is raised in case of error event"""
@@ -177,6 +191,7 @@ class TestPubSubPullSensor:
project_id=TEST_PROJECT,
subscription=TEST_SUBSCRIPTION,
deferrable=True,
+ return_immediately=True,
)
with pytest.raises(AirflowException):
@@ -192,6 +207,7 @@ class TestPubSubPullSensor:
project_id=TEST_PROJECT,
subscription=TEST_SUBSCRIPTION,
deferrable=True,
+ return_immediately=True,
)
test_message = "test"
@@ -238,6 +254,7 @@ class TestPubSubPullSensor:
subscription=TEST_SUBSCRIPTION,
deferrable=True,
messages_callback=messages_callback,
+ return_immediately=True,
)
mock_hook.return_value.pull.return_value = received_messages
@@ -245,3 +262,23 @@ class TestPubSubPullSensor:
resp = operator.execute_complete(context={}, event={"status":
"success", "message": test_message})
mock_log_info.assert_called_with("Sensor pulls messages: %s",
test_message)
assert resp == messages_callback_return_value
+
+ def test_pubsub_pull_sensor_deprecation_warning(self):
+ with pytest.warns(AirflowProviderDeprecationWarning,
match="return_immediately"):
+ sensor = PubSubPullSensor(
+ task_id=TASK_ID,
+ project_id=TEST_PROJECT,
+ subscription=TEST_SUBSCRIPTION,
+ )
+ assert sensor.return_immediately is True
+
+ @pytest.mark.parametrize("return_immediately", [True, False])
+ def test_pubsub_pull_sensor_no_deprecation_warning_when_explicit(self,
return_immediately):
+ with warnings.catch_warnings():
+ warnings.simplefilter("error", AirflowProviderDeprecationWarning)
+ PubSubPullSensor(
+ task_id=TASK_ID,
+ project_id=TEST_PROJECT,
+ subscription=TEST_SUBSCRIPTION,
+ return_immediately=return_immediately,
+ )
diff --git a/providers/google/tests/unit/google/cloud/triggers/test_pubsub.py
b/providers/google/tests/unit/google/cloud/triggers/test_pubsub.py
index 7fb95260b17..bdc28f66276 100644
--- a/providers/google/tests/unit/google/cloud/triggers/test_pubsub.py
+++ b/providers/google/tests/unit/google/cloud/triggers/test_pubsub.py
@@ -16,12 +16,14 @@
# under the License.
from __future__ import annotations
+import warnings
from unittest import mock
import pytest
from google.api_core.exceptions import GoogleAPICallError
from google.cloud.pubsub_v1.types import ReceivedMessage
+from airflow.exceptions import AirflowProviderDeprecationWarning
from airflow.providers.google.cloud.triggers.pubsub import PubsubPullTrigger
from airflow.triggers.base import TriggerEvent
@@ -42,6 +44,7 @@ def trigger():
poke_interval=TEST_POLL_INTERVAL,
gcp_conn_id=TEST_GCP_CONN_ID,
impersonation_chain=None,
+ return_immediately=True,
)
@@ -74,8 +77,34 @@ class TestPubsubPullTrigger:
"poke_interval": TEST_POLL_INTERVAL,
"gcp_conn_id": TEST_GCP_CONN_ID,
"impersonation_chain": None,
+ "return_immediately": True,
}
+ @pytest.mark.asyncio
+
@mock.patch("airflow.providers.google.cloud.hooks.pubsub.PubSubAsyncHook.pull")
+ async def
test_async_pubsub_pull_trigger_passes_return_immediately_false(self, mock_pull):
+ """Test that return_immediately is passed to the hook."""
+ mock_pull.return_value = generate_messages(1)
+ trigger = PubsubPullTrigger(
+ project_id=PROJECT_ID,
+ subscription="subscription",
+ max_messages=MAX_MESSAGES,
+ ack_messages=False,
+ poke_interval=TEST_POLL_INTERVAL,
+ gcp_conn_id=TEST_GCP_CONN_ID,
+ impersonation_chain=None,
+ return_immediately=False,
+ )
+
+ await trigger.run().asend(None)
+
+ mock_pull.assert_called_once_with(
+ project_id=PROJECT_ID,
+ subscription="subscription",
+ max_messages=MAX_MESSAGES,
+ return_immediately=False,
+ )
+
@pytest.mark.asyncio
@mock.patch("airflow.providers.google.cloud.hooks.pubsub.PubSubAsyncHook.pull")
async def test_async_pubsub_pull_trigger_return_event(self, mock_pull):
@@ -88,6 +117,7 @@ class TestPubsubPullTrigger:
poke_interval=TEST_POLL_INTERVAL,
gcp_conn_id=TEST_GCP_CONN_ID,
impersonation_chain=None,
+ return_immediately=True,
)
expected_event = TriggerEvent(
@@ -122,6 +152,7 @@ class TestPubsubPullTrigger:
poke_interval=TEST_POLL_INTERVAL,
gcp_conn_id=TEST_GCP_CONN_ID,
impersonation_chain=None,
+ return_immediately=True,
)
async_hook_actual = trigger.hook
@@ -146,6 +177,7 @@ class TestPubsubPullTrigger:
poke_interval=TEST_POLL_INTERVAL,
gcp_conn_id=TEST_GCP_CONN_ID,
impersonation_chain=None,
+ return_immediately=True,
)
with pytest.raises(GoogleAPICallError, match="Connection error"):
@@ -168,7 +200,41 @@ class TestPubsubPullTrigger:
poke_interval=TEST_POLL_INTERVAL,
gcp_conn_id=TEST_GCP_CONN_ID,
impersonation_chain=None,
+ return_immediately=True,
)
with pytest.raises(GoogleAPICallError, match="Acknowledgement failed"):
await trigger.run().asend(None)
+
+ def test_pubsub_pull_trigger_deprecation_warning(self):
+ test_subscription =
"projects/test_project_id/subscriptions/watcher-subscription"
+ with pytest.warns(AirflowProviderDeprecationWarning,
match="return_immediately") as record:
+ trigger = PubsubPullTrigger(
+ project_id=PROJECT_ID,
+ subscription=test_subscription,
+ max_messages=MAX_MESSAGES,
+ ack_messages=ACK_MESSAGES,
+ poke_interval=TEST_POLL_INTERVAL,
+ gcp_conn_id=TEST_GCP_CONN_ID,
+ impersonation_chain=None,
+ )
+ assert trigger.return_immediately is True
+ # This path is reached from providers/common/messaging's
MessageQueueTrigger.serialize(),
+ # so the warning must name the subscription -- stacklevel=2 otherwise
points at that
+ # unrelated provider's file, leaving the reader no way to tell which
watcher to fix.
+ assert test_subscription in str(record[0].message)
+
+ @pytest.mark.parametrize("return_immediately", [True, False])
+ def test_pubsub_pull_trigger_no_deprecation_warning_when_explicit(self,
return_immediately):
+ with warnings.catch_warnings():
+ warnings.simplefilter("error", AirflowProviderDeprecationWarning)
+ PubsubPullTrigger(
+ project_id=PROJECT_ID,
+ subscription="subscription",
+ max_messages=MAX_MESSAGES,
+ ack_messages=ACK_MESSAGES,
+ poke_interval=TEST_POLL_INTERVAL,
+ gcp_conn_id=TEST_GCP_CONN_ID,
+ impersonation_chain=None,
+ return_immediately=return_immediately,
+ )