From 45112f91567106f717a2088edabf94b24a5f0eb5 Mon Sep 17 00:00:00 2001 From: Mike Prieto Date: Mon, 21 Sep 2026 23:12:46 +0300 Subject: [PATCH 1/4] 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. --- .../google/cloud/operators/pubsub.py | 22 +++++++- .../providers/google/cloud/sensors/pubsub.py | 15 +++++- .../providers/google/cloud/triggers/pubsub.py | 22 +++++++- .../google/cloud/operators/test_pubsub.py | 53 ++++++++++++++++++- .../unit/google/cloud/sensors/test_pubsub.py | 51 ++++++++++++++++++ .../unit/google/cloud/triggers/test_pubsub.py | 42 +++++++++++++++ 6 files changed, 200 insertions(+), 5 deletions(-) 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 33a1b302786cc..207ca834004b0 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 @@ from __future__ import annotations +import warnings from collections.abc import Callable, Sequence from functools import cached_property from typing import TYPE_CHECKING, Any @@ -42,6 +43,7 @@ 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 @@ -799,6 +801,13 @@ 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: 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. """ template_fields: Sequence[str] = ( @@ -819,6 +828,7 @@ def __init__( 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 +841,15 @@ def __init__( self.impersonation_chain = impersonation_chain self.deferrable = deferrable self.poll_interval = poll_interval + if return_immediately is not None: + warnings.warn( + "The default value of `return_immediately` will be changed to `False` in a future major release.", + AirflowProviderDeprecationWarning, + stacklevel=2, + ) + self.return_immediately = return_immediately + else: + self.return_immediately = True def execute(self, context: Context) -> list: if self.deferrable: @@ -843,6 +862,7 @@ def execute(self, context: Context) -> list: 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 +875,7 @@ def execute(self, context: Context) -> list: 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 f138271b66e4c..7c9ceccd2774a 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,6 +27,7 @@ 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 @@ -113,7 +115,7 @@ def __init__( 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,13 +129,21 @@ def __init__( self.project_id = project_id self.subscription = subscription self.max_messages = max_messages - self.return_immediately = return_immediately self.ack_messages = ack_messages self.messages_callback = messages_callback self.impersonation_chain = impersonation_chain self.deferrable = deferrable self.poke_interval = poke_interval self._return_value = None + if return_immediately is not None: + warnings.warn( + "The default value of `return_immediately` will be changed to `False` in a future major release.", + AirflowProviderDeprecationWarning, + stacklevel=2, + ) + self.return_immediately = return_immediately + else: + self.return_immediately = True def poke(self, context: Context) -> bool: hook = PubSubHook( @@ -176,6 +186,7 @@ def execute(self, context: Context) -> None: 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 22314838c4649..4c8aff5d27ac4 100644 --- a/providers/google/src/airflow/providers/google/cloud/triggers/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/triggers/pubsub.py @@ -19,12 +19,14 @@ 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.version_compat import AIRFLOW_V_3_0_PLUS from airflow.triggers.base import TriggerEvent @@ -55,6 +57,13 @@ 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: 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. """ def __init__( @@ -66,6 +75,7 @@ def __init__( 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 +85,15 @@ def __init__( self.poke_interval = poke_interval self.gcp_conn_id = gcp_conn_id self.impersonation_chain = impersonation_chain + if return_immediately is not None: + warnings.warn( + "The default value of `return_immediately` will be changed to `False` in a future major release.", + AirflowProviderDeprecationWarning, + stacklevel=2, + ) + self.return_immediately = return_immediately + else: + self.return_immediately = True def serialize(self) -> tuple[str, dict[str, Any]]: """Serialize PubsubPullTrigger arguments and classpath.""" @@ -88,6 +107,7 @@ def serialize(self) -> tuple[str, dict[str, Any]]: "poke_interval": self.poke_interval, "gcp_conn_id": self.gcp_conn_id, "impersonation_chain": self.impersonation_chain, + "return_immediately": self.return_immediately, }, ) @@ -97,7 +117,7 @@ async def run(self) -> AsyncIterator[TriggerEvent]: 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/tests/unit/google/cloud/operators/test_pubsub.py b/providers/google/tests/unit/google/cloud/operators/test_pubsub.py index 3537c5266db2e..f8570d14dd253 100644 --- a/providers/google/tests/unit/google/cloud/operators/test_pubsub.py +++ b/providers/google/tests/unit/google/cloud/operators/test_pubsub.py @@ -25,6 +25,7 @@ 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 +35,9 @@ PubSubPublishMessageOperator, PubSubPullOperator, ) +from airflow.providers.google.cloud.triggers.pubsub import PubsubPullTrigger + +pytestmark = pytest.mark.filterwarnings("ignore::airflow.exceptions.AirflowProviderDeprecationWarning") TASK_ID = "test-task-id" TEST_PROJECT = "test-project" @@ -513,6 +517,21 @@ def messages_callback( assert response == messages_callback_return_value + @mock.patch("airflow.providers.google.cloud.operators.pubsub.PubSubHook") + def test_execute_with_return_immediately_false(self, mock_hook): + operator = PubSubPullOperator( + task_id=TASK_ID, + project_id=TEST_PROJECT, + subscription=TEST_SUBSCRIPTION, + return_immediately=False, + ) + + mock_hook.return_value.pull.return_value = [] + operator.execute({}) + mock_hook.return_value.pull.assert_called_once_with( + project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, max_messages=5, return_immediately=False + ) + @pytest.mark.db_test @mock.patch("airflow.providers.google.cloud.operators.pubsub.PubSubHook") def test_execute_deferred(self, mock_hook): @@ -526,9 +545,32 @@ def test_execute_deferred(self, mock_hook): subscription=TEST_SUBSCRIPTION, deferrable=True, ) - 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 True + + @pytest.mark.db_test + @mock.patch("airflow.providers.google.cloud.operators.pubsub.PubSubHook") + def test_execute_deferred_with_return_immediately_false(self, mock_hook): + """ + Asserts that a task is deferred and a PubSubPullOperator will be fired + when the PubSubPullOperator is executed with deferrable=True. + """ + task = PubSubPullOperator( + task_id=TASK_ID, + project_id=TEST_PROJECT, + subscription=TEST_SUBSCRIPTION, + return_immediately=False, + deferrable=True, + ) + with pytest.raises(TaskDeferred) as exc: task.execute(mock.MagicMock()) + assert isinstance(exc.value.trigger, PubsubPullTrigger) + assert exc.value.trigger.return_immediately is False + @mock.patch("airflow.providers.google.cloud.operators.pubsub.PubSubHook") def test_get_openlineage_facets(self, mock_hook): operator = PubSubPullOperator( @@ -631,3 +673,12 @@ def test_execute_complete_use_default_message_callback(self, mock_hook): 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"): + PubSubPullOperator( + task_id=TASK_ID, + project_id=TEST_PROJECT, + subscription=TEST_SUBSCRIPTION, + return_immediately=False, + ) 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 4cd1b48fbfb60..64a6d72dd511a 100644 --- a/providers/google/tests/unit/google/cloud/sensors/test_pubsub.py +++ b/providers/google/tests/unit/google/cloud/sensors/test_pubsub.py @@ -24,10 +24,13 @@ 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 +pytestmark = pytest.mark.filterwarnings("ignore::airflow.exceptions.AirflowProviderDeprecationWarning") + TASK_ID = "test-task-id" TEST_PROJECT = "test-project" TEST_SUBSCRIPTION = "test-subscription" @@ -99,6 +102,26 @@ def test_execute(self, mock_hook): ) assert generated_dicts == response + @mock.patch("airflow.providers.google.cloud.sensors.pubsub.PubSubHook") + def test_execute_with_return_immediately_false(self, mock_hook): + operator = PubSubPullSensor( + task_id=TASK_ID, + project_id=TEST_PROJECT, + subscription=TEST_SUBSCRIPTION, + poke_interval=0, + return_immediately=False, + ) + + generated_messages = self._generate_messages(5) + generated_dicts = self._generate_dicts(5) + mock_hook.return_value.pull.return_value = generated_messages + + response = operator.execute({}) + mock_hook.return_value.pull.assert_called_once_with( + project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, max_messages=5, return_immediately=False + ) + assert generated_dicts == response + @mock.patch("airflow.providers.google.cloud.sensors.pubsub.PubSubHook") def test_execute_timeout(self, mock_hook): operator = PubSubPullSensor( @@ -167,6 +190,25 @@ def test_pubsub_pull_sensor_async(self): 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 True + + def test_pubsub_pull_sensor_async_with_return_immediately_false(self): + """ + Asserts that a task is deferred and a PubsubPullTrigger will be fired + with custom return_immediately value. + """ + task = PubSubPullSensor( + task_id="test_task_id", + ack_messages=True, + project_id=TEST_PROJECT, + subscription=TEST_SUBSCRIPTION, + deferrable=True, + return_immediately=False, + ) + 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 False def test_pubsub_pull_sensor_async_execute_should_throw_exception(self): """Tests that an AirflowException is raised in case of error event""" @@ -245,3 +287,12 @@ def messages_callback( 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"): + PubSubPullSensor( + task_id=TASK_ID, + project_id=TEST_PROJECT, + subscription=TEST_SUBSCRIPTION, + return_immediately=False, + ) 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 7fb95260b17fa..9ff897be74bb9 100644 --- a/providers/google/tests/unit/google/cloud/triggers/test_pubsub.py +++ b/providers/google/tests/unit/google/cloud/triggers/test_pubsub.py @@ -22,9 +22,12 @@ 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 +pytestmark = pytest.mark.filterwarnings("ignore::airflow.exceptions.AirflowProviderDeprecationWarning") + TEST_POLL_INTERVAL = 10 TEST_GCP_CONN_ID = "google_cloud_default" PROJECT_ID = "test_project_id" @@ -74,8 +77,34 @@ def test_async_pubsub_pull_trigger_serialization_should_execute_successfully(sel "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): @@ -172,3 +201,16 @@ async def test_async_pubsub_pull_trigger_exception_during_ack(self, mock_pull, m with pytest.raises(GoogleAPICallError, match="Acknowledgement failed"): await trigger.run().asend(None) + + def test_pubsub_pull_trigger_deprecation_warning(self): + with pytest.warns(AirflowProviderDeprecationWarning, match="return_immediately"): + 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=False, + ) From 965b3bbe12308330e32ee539a2dd01ee2f0887cc Mon Sep 17 00:00:00 2001 From: Shahar Epstein <60007259+shahar1@users.noreply.github.com> Date: Mon, 21 Sep 2026 23:19:15 +0300 Subject: [PATCH 2/4] 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. --- .../google/cloud/operators/pubsub.py | 40 ++++++---- .../providers/google/cloud/sensors/pubsub.py | 38 +++++----- .../providers/google/cloud/triggers/pubsub.py | 26 ++++--- .../airflow/providers/google/common/consts.py | 8 ++ .../google/cloud/operators/test_pubsub.py | 72 ++++++++---------- .../unit/google/cloud/sensors/test_pubsub.py | 74 ++++++++----------- .../unit/google/cloud/triggers/test_pubsub.py | 32 +++++++- 7 files changed, 154 insertions(+), 136 deletions(-) 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 207ca834004b0..6caad672cdab7 100644 --- a/providers/google/src/airflow/providers/google/cloud/operators/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/operators/pubsub.py @@ -49,7 +49,10 @@ 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: @@ -758,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: @@ -801,13 +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: 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. """ template_fields: Sequence[str] = ( @@ -841,15 +850,14 @@ def __init__( self.impersonation_chain = impersonation_chain self.deferrable = deferrable self.poll_interval = poll_interval - if return_immediately is not None: + if return_immediately is None: warnings.warn( - "The default value of `return_immediately` will be changed to `False` in a future major release.", + PUBSUB_RETURN_IMMEDIATELY_DEPRECATION_MESSAGE, AirflowProviderDeprecationWarning, stacklevel=2, ) - self.return_immediately = return_immediately - else: - self.return_immediately = True + return_immediately = True + self.return_immediately = return_immediately def execute(self, context: Context) -> list: if self.deferrable: 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 7c9ceccd2774a..5e993d5be4632 100644 --- a/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py @@ -31,6 +31,7 @@ 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 @@ -51,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 @@ -63,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. @@ -75,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 @@ -129,21 +130,20 @@ def __init__( 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 self.impersonation_chain = impersonation_chain self.deferrable = deferrable self.poke_interval = poke_interval self._return_value = None - if return_immediately is not None: - warnings.warn( - "The default value of `return_immediately` will be changed to `False` in a future major release.", - AirflowProviderDeprecationWarning, - stacklevel=2, - ) - self.return_immediately = return_immediately - else: - self.return_immediately = True def poke(self, context: Context) -> bool: hook = PubSubHook( 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 4c8aff5d27ac4..eabe74055cfa8 100644 --- a/providers/google/src/airflow/providers/google/cloud/triggers/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/triggers/pubsub.py @@ -28,6 +28,7 @@ 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 @@ -57,13 +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: 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: 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__( @@ -85,15 +88,14 @@ def __init__( self.poke_interval = poke_interval self.gcp_conn_id = gcp_conn_id self.impersonation_chain = impersonation_chain - if return_immediately is not None: + if return_immediately is None: warnings.warn( - "The default value of `return_immediately` will be changed to `False` in a future major release.", + f"{PUBSUB_RETURN_IMMEDIATELY_DEPRECATION_MESSAGE} Subscription: {self.subscription}.", AirflowProviderDeprecationWarning, stacklevel=2, ) - self.return_immediately = return_immediately - else: - self.return_immediately = True + return_immediately = True + self.return_immediately = return_immediately def serialize(self) -> tuple[str, dict[str, Any]]: """Serialize PubsubPullTrigger arguments and classpath.""" diff --git a/providers/google/src/airflow/providers/google/common/consts.py b/providers/google/src/airflow/providers/google/common/consts.py index f8d7209901d73..cd7ae3d912b77 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 @@ 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/tests/unit/google/cloud/operators/test_pubsub.py b/providers/google/tests/unit/google/cloud/operators/test_pubsub.py index f8570d14dd253..946a79d6a81f6 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 @@ -37,8 +38,6 @@ ) from airflow.providers.google.cloud.triggers.pubsub import PubsubPullTrigger -pytestmark = pytest.mark.filterwarnings("ignore::airflow.exceptions.AirflowProviderDeprecationWarning") - TASK_ID = "test-task-id" TEST_PROJECT = "test-project" TEST_TOPIC = "test-topic" @@ -449,16 +448,24 @@ def _generate_messages(self, count): 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): @@ -467,6 +474,7 @@ def test_execute_with_ack_messages(self, mock_hook): project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, ack_messages=True, + return_immediately=True, ) generated_messages = self._generate_messages(5) @@ -504,6 +512,7 @@ def messages_callback( project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, messages_callback=messages_callback, + return_immediately=True, ) mock_hook.return_value.pull.return_value = generated_messages @@ -517,43 +526,10 @@ def messages_callback( assert response == messages_callback_return_value - @mock.patch("airflow.providers.google.cloud.operators.pubsub.PubSubHook") - def test_execute_with_return_immediately_false(self, mock_hook): - operator = PubSubPullOperator( - task_id=TASK_ID, - project_id=TEST_PROJECT, - subscription=TEST_SUBSCRIPTION, - return_immediately=False, - ) - - mock_hook.return_value.pull.return_value = [] - operator.execute({}) - mock_hook.return_value.pull.assert_called_once_with( - project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, max_messages=5, return_immediately=False - ) - - @pytest.mark.db_test - @mock.patch("airflow.providers.google.cloud.operators.pubsub.PubSubHook") - def test_execute_deferred(self, mock_hook): - """ - Asserts that a task is deferred and a PubSubPullOperator will be fired - when the PubSubPullOperator is executed with deferrable=True. - """ - task = PubSubPullOperator( - task_id=TASK_ID, - project_id=TEST_PROJECT, - subscription=TEST_SUBSCRIPTION, - deferrable=True, - ) - with pytest.raises(TaskDeferred) as exc: - task.execute(mock.MagicMock()) - - assert isinstance(exc.value.trigger, PubsubPullTrigger) - assert exc.value.trigger.return_immediately is True - @pytest.mark.db_test + @pytest.mark.parametrize("return_immediately", [True, False]) @mock.patch("airflow.providers.google.cloud.operators.pubsub.PubSubHook") - def test_execute_deferred_with_return_immediately_false(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. @@ -562,14 +538,14 @@ def test_execute_deferred_with_return_immediately_false(self, mock_hook): task_id=TASK_ID, project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, - return_immediately=False, deferrable=True, + return_immediately=return_immediately, ) with pytest.raises(TaskDeferred) as exc: task.execute(mock.MagicMock()) assert isinstance(exc.value.trigger, PubsubPullTrigger) - assert exc.value.trigger.return_immediately is False + 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): @@ -577,6 +553,7 @@ def test_get_openlineage_facets(self, mock_hook): task_id=TASK_ID, project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, + return_immediately=True, ) generated_messages = self._generate_messages(5) @@ -635,6 +612,7 @@ def messages_callback( subscription=TEST_SUBSCRIPTION, deferrable=True, messages_callback=messages_callback, + return_immediately=True, ) mock_hook.return_value.pull.return_value = received_messages @@ -666,6 +644,7 @@ def test_execute_complete_use_default_message_callback(self, mock_hook): project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, deferrable=True, + return_immediately=True, ) mock_hook.return_value.pull.return_value = received_messages @@ -676,9 +655,20 @@ def test_execute_complete_use_default_message_callback(self, mock_hook): 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=False, + 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 64a6d72dd511a..d5f4e692b0f13 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 @@ -29,8 +30,6 @@ from airflow.providers.google.cloud.sensors.pubsub import PubSubPullSensor from airflow.providers.google.cloud.triggers.pubsub import PubsubPullTrigger -pytestmark = pytest.mark.filterwarnings("ignore::airflow.exceptions.AirflowProviderDeprecationWarning") - TASK_ID = "test-task-id" TEST_PROJECT = "test-project" TEST_SUBSCRIPTION = "test-subscription" @@ -58,6 +57,7 @@ def test_poke_no_messages(self, mock_hook): task_id=TASK_ID, project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, + return_immediately=True, ) mock_hook.return_value.pull.return_value = [] @@ -70,6 +70,7 @@ def test_poke_with_ack_messages(self, mock_hook): project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, ack_messages=True, + return_immediately=True, ) generated_messages = self._generate_messages(5) @@ -83,13 +84,15 @@ def test_poke_with_ack_messages(self, mock_hook): 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) @@ -98,27 +101,10 @@ def test_execute(self, mock_hook): 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 - ) - assert generated_dicts == response - - @mock.patch("airflow.providers.google.cloud.sensors.pubsub.PubSubHook") - def test_execute_with_return_immediately_false(self, mock_hook): - operator = PubSubPullSensor( - task_id=TASK_ID, project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, - poke_interval=0, - return_immediately=False, - ) - - generated_messages = self._generate_messages(5) - generated_dicts = self._generate_dicts(5) - mock_hook.return_value.pull.return_value = generated_messages - - response = operator.execute({}) - mock_hook.return_value.pull.assert_called_once_with( - project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, max_messages=5, return_immediately=False + max_messages=5, + return_immediately=return_immediately, ) assert generated_dicts == response @@ -130,6 +116,7 @@ def test_execute_timeout(self, mock_hook): subscription=TEST_SUBSCRIPTION, poke_interval=0, timeout=1, + return_immediately=True, ) mock_hook.return_value.pull.return_value = [] @@ -162,6 +149,7 @@ def messages_callback( subscription=TEST_SUBSCRIPTION, poke_interval=0, messages_callback=messages_callback, + return_immediately=True, ) mock_hook.return_value.pull.return_value = generated_messages @@ -175,27 +163,11 @@ def messages_callback( assert response == messages_callback_return_value - def test_pubsub_pull_sensor_async(self): - """ - Asserts that a task is deferred and a PubsubPullTrigger will be fired - when the PubSubPullSensor is executed. - """ - task = PubSubPullSensor( - task_id="test_task_id", - ack_messages=True, - project_id=TEST_PROJECT, - subscription=TEST_SUBSCRIPTION, - deferrable=True, - ) - 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 True - - def test_pubsub_pull_sensor_async_with_return_immediately_false(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 - with custom return_immediately value. + when the PubSubPullSensor is executed, with the configured return_immediately value. """ task = PubSubPullSensor( task_id="test_task_id", @@ -203,12 +175,12 @@ def test_pubsub_pull_sensor_async_with_return_immediately_false(self): project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, deferrable=True, - return_immediately=False, + 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 False + 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""" @@ -219,6 +191,7 @@ def test_pubsub_pull_sensor_async_execute_should_throw_exception(self): project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, deferrable=True, + return_immediately=True, ) with pytest.raises(AirflowException): @@ -234,6 +207,7 @@ def test_pubsub_pull_sensor_async_execute_complete(self): project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, deferrable=True, + return_immediately=True, ) test_message = "test" @@ -280,6 +254,7 @@ def messages_callback( subscription=TEST_SUBSCRIPTION, deferrable=True, messages_callback=messages_callback, + return_immediately=True, ) mock_hook.return_value.pull.return_value = received_messages @@ -290,9 +265,20 @@ def messages_callback( 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=False, + 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 9ff897be74bb9..bdc28f66276a2 100644 --- a/providers/google/tests/unit/google/cloud/triggers/test_pubsub.py +++ b/providers/google/tests/unit/google/cloud/triggers/test_pubsub.py @@ -16,6 +16,7 @@ # under the License. from __future__ import annotations +import warnings from unittest import mock import pytest @@ -26,8 +27,6 @@ from airflow.providers.google.cloud.triggers.pubsub import PubsubPullTrigger from airflow.triggers.base import TriggerEvent -pytestmark = pytest.mark.filterwarnings("ignore::airflow.exceptions.AirflowProviderDeprecationWarning") - TEST_POLL_INTERVAL = 10 TEST_GCP_CONN_ID = "google_cloud_default" PROJECT_ID = "test_project_id" @@ -45,6 +44,7 @@ def trigger(): poke_interval=TEST_POLL_INTERVAL, gcp_conn_id=TEST_GCP_CONN_ID, impersonation_chain=None, + return_immediately=True, ) @@ -117,6 +117,7 @@ async def test_async_pubsub_pull_trigger_return_event(self, mock_pull): poke_interval=TEST_POLL_INTERVAL, gcp_conn_id=TEST_GCP_CONN_ID, impersonation_chain=None, + return_immediately=True, ) expected_event = TriggerEvent( @@ -151,6 +152,7 @@ def test_hook(self, mock_async_hook): poke_interval=TEST_POLL_INTERVAL, gcp_conn_id=TEST_GCP_CONN_ID, impersonation_chain=None, + return_immediately=True, ) async_hook_actual = trigger.hook @@ -175,6 +177,7 @@ async def test_async_pubsub_pull_trigger_exception_during_pull(self, mock_pull): 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"): @@ -197,13 +200,34 @@ async def test_async_pubsub_pull_trigger_exception_during_ack(self, mock_pull, m 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): - with pytest.warns(AirflowProviderDeprecationWarning, match="return_immediately"): + 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", @@ -212,5 +236,5 @@ def test_pubsub_pull_trigger_deprecation_warning(self): poke_interval=TEST_POLL_INTERVAL, gcp_conn_id=TEST_GCP_CONN_ID, impersonation_chain=None, - return_immediately=False, + return_immediately=return_immediately, ) From 69f179cc8fba99ee1cb8645b0f76e34d63fa3da6 Mon Sep 17 00:00:00 2001 From: Shahar Epstein <60007259+shahar1@users.noreply.github.com> Date: Mon, 21 Sep 2026 23:20:26 +0300 Subject: [PATCH 3/4] 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. --- .../providers/google/event_scheduling/events/pubsub.py | 1 + .../tests/system/google/cloud/pubsub/example_pubsub.py | 5 +++++ .../system/google/cloud/pubsub/example_pubsub_deferrable.py | 1 + .../google/event_scheduling/example_event_schedule_pubsub.py | 1 + 4 files changed, 8 insertions(+) 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 1d1e361218151..3c808e9fc6b56 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 f2b6dae0158cc..de6608aad268b 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 @@ 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 @@ # [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 195f5a313a51f..78c92f44c5fa8 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 @@ 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 4dbc466cc7d2e..93a2aa581e7b9 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 @@ 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 From e630f54df6b0a8d6bf0ea810806c347557d1b657 Mon Sep 17 00:00:00 2001 From: Shahar Epstein <60007259+shahar1@users.noreply.github.com> Date: Mon, 21 Sep 2026 23:20:43 +0300 Subject: [PATCH 4/4] 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. --- providers/google/docs/changelog.rst | 14 ++++++++++++++ providers/google/docs/operators/cloud/pubsub.rst | 7 +++++++ 2 files changed, 21 insertions(+) diff --git a/providers/google/docs/changelog.rst b/providers/google/docs/changelog.rst index 5bb751bb5e3f8..86ac6b0cf7ed1 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 a113b2e775f73..2376dcbe29108 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]