From 3ec7a55c99edb2a7b37af17c53aac76906965c04 Mon Sep 17 00:00:00 2001 From: Mike Prieto Date: Wed, 27 May 2026 21:38:12 +0000 Subject: [PATCH 1/7] Make return_immediately default to False and be configurable for all Pub/Sub modules --- .../google/cloud/operators/pubsub.py | 12 ++++++++- .../providers/google/cloud/sensors/pubsub.py | 3 ++- .../providers/google/cloud/triggers/pubsub.py | 12 ++++++++- .../google/cloud/operators/test_pubsub.py | 26 ++++++++++++++++--- .../unit/google/cloud/sensors/test_pubsub.py | 19 ++++++++++++++ .../unit/google/cloud/triggers/test_pubsub.py | 26 +++++++++++++++++++ 6 files changed, 92 insertions(+), 6 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..244aeb6260578 100644 --- a/providers/google/src/airflow/providers/google/cloud/operators/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/operators/pubsub.py @@ -799,6 +799,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 +826,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 = False, **kwargs, ) -> None: super().__init__(**kwargs) @@ -831,6 +839,7 @@ def __init__( self.impersonation_chain = impersonation_chain self.deferrable = deferrable self.poll_interval = poll_interval + self.return_immediately = return_immediately def execute(self, context: Context) -> list: if self.deferrable: @@ -843,6 +852,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 +865,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..007b2c92081fe 100644 --- a/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py @@ -113,7 +113,7 @@ def __init__( project_id: str, subscription: str, max_messages: int = 5, - return_immediately: bool = True, + return_immediately: bool = False, ack_messages: bool = False, gcp_conn_id: str = "google_cloud_default", messages_callback: Callable[[list[ReceivedMessage], Context], Any] | None = None, @@ -176,6 +176,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..b1ea539a250d8 100644 --- a/providers/google/src/airflow/providers/google/cloud/triggers/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/triggers/pubsub.py @@ -55,6 +55,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 +73,7 @@ def __init__( gcp_conn_id: str, poke_interval: float = 10.0, impersonation_chain: str | Sequence[str] | None = None, + return_immediately: bool = False, ): super().__init__() self.project_id = project_id @@ -75,6 +83,7 @@ def __init__( self.poke_interval = poke_interval self.gcp_conn_id = gcp_conn_id self.impersonation_chain = impersonation_chain + self.return_immediately = return_immediately def serialize(self) -> tuple[str, dict[str, Any]]: """Serialize PubsubPullTrigger arguments and classpath.""" @@ -88,6 +97,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 +107,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..43c0a6abb8488 100644 --- a/providers/google/tests/unit/google/cloud/operators/test_pubsub.py +++ b/providers/google/tests/unit/google/cloud/operators/test_pubsub.py @@ -34,6 +34,7 @@ PubSubPublishMessageOperator, PubSubPullOperator, ) +from airflow.providers.google.cloud.triggers.pubsub import PubsubPullTrigger TASK_ID = "test-task-id" TEST_PROJECT = "test-project" @@ -506,13 +507,28 @@ def messages_callback( 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=False ) messages_callback.assert_called_once() assert response == messages_callback_return_value + @mock.patch("airflow.providers.google.cloud.operators.pubsub.PubSubHook") + def test_execute_with_return_immediately_true(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=True + ) + @pytest.mark.db_test @mock.patch("airflow.providers.google.cloud.operators.pubsub.PubSubHook") def test_execute_deferred(self, mock_hook): @@ -525,10 +541,14 @@ def test_execute_deferred(self, mock_hook): project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, deferrable=True, + return_immediately=False, ) - 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 False + @mock.patch("airflow.providers.google.cloud.operators.pubsub.PubSubHook") def test_get_openlineage_facets(self, mock_hook): operator = PubSubPullOperator( @@ -543,7 +563,7 @@ def test_get_openlineage_facets(self, mock_hook): assert generated_dicts == 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=False ) result = operator.get_openlineage_facets_on_complete(operator) 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..be339e8fd2c1f 100644 --- a/providers/google/tests/unit/google/cloud/sensors/test_pubsub.py +++ b/providers/google/tests/unit/google/cloud/sensors/test_pubsub.py @@ -167,6 +167,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 False + + def test_pubsub_pull_sensor_async_with_return_immediately_true(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=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_execute_should_throw_exception(self): """Tests that an AirflowException is raised in case of error event""" 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..b18bfc2246d4f 100644 --- a/providers/google/tests/unit/google/cloud/triggers/test_pubsub.py +++ b/providers/google/tests/unit/google/cloud/triggers/test_pubsub.py @@ -74,8 +74,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": False, } + @pytest.mark.asyncio + @mock.patch("airflow.providers.google.cloud.hooks.pubsub.PubSubAsyncHook.pull") + async def test_async_pubsub_pull_trigger_passes_return_immediately(self, mock_pull): + """Test that return_immediately is passed to the hook.""" + mock_pull.return_value = await 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=True, + ) + + await trigger.run().asend(None) + + mock_pull.assert_called_once_with( + project_id=PROJECT_ID, + subscription="subscription", + max_messages=MAX_MESSAGES, + return_immediately=True, + ) + @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): From 3b5e6dd7b70478f4c646b09e82d85968915d88fb Mon Sep 17 00:00:00 2001 From: Mike Prieto Date: Tue, 2 Jun 2026 16:36:59 +0000 Subject: [PATCH 2/7] Update test expectations for having return_immediately default to false for Pub/Sub modules --- .../google/tests/unit/google/cloud/operators/test_pubsub.py | 2 +- .../google/tests/unit/google/cloud/sensors/test_pubsub.py | 4 ++-- .../google/tests/unit/google/cloud/triggers/test_pubsub.py | 2 +- 3 files changed, 4 insertions(+), 4 deletions(-) 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 43c0a6abb8488..2058e93a49dd9 100644 --- a/providers/google/tests/unit/google/cloud/operators/test_pubsub.py +++ b/providers/google/tests/unit/google/cloud/operators/test_pubsub.py @@ -520,7 +520,7 @@ def test_execute_with_return_immediately_true(self, mock_hook): task_id=TASK_ID, project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, - return_immediately=False, + return_immediately=True, ) mock_hook.return_value.pull.return_value = [] 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 be339e8fd2c1f..cb73498d6a4be 100644 --- a/providers/google/tests/unit/google/cloud/sensors/test_pubsub.py +++ b/providers/google/tests/unit/google/cloud/sensors/test_pubsub.py @@ -95,7 +95,7 @@ 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 + project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, max_messages=5, return_immediately=False ) assert generated_dicts == response @@ -145,7 +145,7 @@ def messages_callback( 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=False ) messages_callback.assert_called_once() 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 b18bfc2246d4f..f2c0dd4ef7cf7 100644 --- a/providers/google/tests/unit/google/cloud/triggers/test_pubsub.py +++ b/providers/google/tests/unit/google/cloud/triggers/test_pubsub.py @@ -81,7 +81,7 @@ def test_async_pubsub_pull_trigger_serialization_should_execute_successfully(sel @mock.patch("airflow.providers.google.cloud.hooks.pubsub.PubSubAsyncHook.pull") async def test_async_pubsub_pull_trigger_passes_return_immediately(self, mock_pull): """Test that return_immediately is passed to the hook.""" - mock_pull.return_value = await generate_messages(1) + mock_pull.return_value = generate_messages(1) trigger = PubsubPullTrigger( project_id=PROJECT_ID, subscription="subscription", From ea7f6c31f2a33f3e66ffab7dbfd321784bad7d6e Mon Sep 17 00:00:00 2001 From: Mike Prieto Date: Mon, 6 Jul 2026 17:43:46 +0000 Subject: [PATCH 3/7] Change return_immediately default value back to True and add deprecation warnings for the setting to Pub/Sub modules --- .../airflow/providers/google/cloud/hooks/pubsub.py | 12 ++++++++++++ .../providers/google/cloud/operators/pubsub.py | 12 +++++++++++- .../providers/google/cloud/sensors/pubsub.py | 12 +++++++++++- .../providers/google/cloud/triggers/pubsub.py | 12 +++++++++++- .../unit/google/cloud/operators/test_pubsub.py | 13 ++++++------- .../tests/unit/google/cloud/sensors/test_pubsub.py | 12 ++++++------ .../tests/unit/google/cloud/triggers/test_pubsub.py | 8 ++++---- 7 files changed, 61 insertions(+), 20 deletions(-) diff --git a/providers/google/src/airflow/providers/google/cloud/hooks/pubsub.py b/providers/google/src/airflow/providers/google/cloud/hooks/pubsub.py index a983fe3e4b7c7..ab1b1cee4c5f7 100644 --- a/providers/google/src/airflow/providers/google/cloud/hooks/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/hooks/pubsub.py @@ -520,6 +520,12 @@ def pull( the base64-encoded message content. See https://cloud.google.com/pubsub/docs/reference/rpc/google.pubsub.v1#google.pubsub.v1.ReceivedMessage """ + warnings.warn( + "The `return_immediately` parameter is deprecated and will be removed in a future release.", + AirflowProviderDeprecationWarning, + stacklevel=2, + ) + subscriber = self.subscriber_client # E501 subscription_path = f"projects/{project_id}/subscriptions/{subscription}" @@ -712,6 +718,12 @@ async def pull( the base64-encoded message content. See https://cloud.google.com/pubsub/docs/reference/rpc/google.pubsub.v1#google.pubsub.v1.ReceivedMessage """ + warnings.warn( + "The `return_immediately` parameter is deprecated and will be removed in a future release.", + AirflowProviderDeprecationWarning, + stacklevel=2, + ) + subscriber = await self._get_subscriber_client() subscription_path = f"projects/{project_id}/subscriptions/{subscription}" self.log.info("Pulling max %d messages from subscription (path) %s", max_messages, subscription_path) 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 244aeb6260578..3cc7c69c19c3f 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 @@ -826,7 +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 = False, + return_immediately: bool = True, **kwargs, ) -> None: super().__init__(**kwargs) @@ -841,6 +843,14 @@ def __init__( self.poll_interval = poll_interval self.return_immediately = return_immediately + warnings.warn( + "The `return_immediately` parameter is deprecated and will be removed in a future release. " + "Its default value will be changed to `False` in the next major release. " + "Planned removal date: August 01, 2026.", + AirflowProviderDeprecationWarning, + stacklevel=2, + ) + def execute(self, context: Context) -> list: if self.deferrable: self.defer( 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 007b2c92081fe..a266f42d76e8a 100644 --- a/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py @@ -22,10 +22,12 @@ from collections.abc import Callable, Sequence from datetime import timedelta from typing import TYPE_CHECKING, Any +import warnings 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 = False, + return_immediately: bool = True, ack_messages: bool = False, gcp_conn_id: str = "google_cloud_default", messages_callback: Callable[[list[ReceivedMessage], Context], Any] | None = None, @@ -135,6 +137,14 @@ def __init__( self.poke_interval = poke_interval self._return_value = None + warnings.warn( + "The `return_immediately` parameter is deprecated and will be removed in a future release. " + "Its default value will be changed to `False` in the next major release. " + "Planned removal date: August 01, 2026.", + AirflowProviderDeprecationWarning, + stacklevel=2, + ) + def poke(self, context: Context) -> bool: hook = PubSubHook( gcp_conn_id=self.gcp_conn_id, 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 b1ea539a250d8..7e87f0c9039df 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 @@ -73,7 +75,7 @@ def __init__( gcp_conn_id: str, poke_interval: float = 10.0, impersonation_chain: str | Sequence[str] | None = None, - return_immediately: bool = False, + return_immediately: bool = True, ): super().__init__() self.project_id = project_id @@ -85,6 +87,14 @@ def __init__( self.impersonation_chain = impersonation_chain self.return_immediately = return_immediately + warnings.warn( + "The `return_immediately` parameter is deprecated and will be removed in a future release. " + "Its default value will be changed to `False` in the next major release. " + "Planned removal date: August 01, 2026.", + AirflowProviderDeprecationWarning, + stacklevel=2, + ) + def serialize(self) -> tuple[str, dict[str, Any]]: """Serialize PubsubPullTrigger arguments and classpath.""" return ( 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 2058e93a49dd9..680aa5d96c694 100644 --- a/providers/google/tests/unit/google/cloud/operators/test_pubsub.py +++ b/providers/google/tests/unit/google/cloud/operators/test_pubsub.py @@ -507,7 +507,7 @@ def messages_callback( 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 + project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, max_messages=5, return_immediately=True ) messages_callback.assert_called_once() @@ -515,18 +515,18 @@ def messages_callback( assert response == messages_callback_return_value @mock.patch("airflow.providers.google.cloud.operators.pubsub.PubSubHook") - def test_execute_with_return_immediately_true(self, mock_hook): + 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=True, + 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=True + project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, max_messages=5, return_immediately=False ) @pytest.mark.db_test @@ -541,13 +541,12 @@ def test_execute_deferred(self, mock_hook): project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, deferrable=True, - return_immediately=False, ) 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 True @mock.patch("airflow.providers.google.cloud.operators.pubsub.PubSubHook") def test_get_openlineage_facets(self, mock_hook): @@ -563,7 +562,7 @@ def test_get_openlineage_facets(self, mock_hook): assert generated_dicts == operator.execute({}) mock_hook.return_value.pull.assert_called_once_with( - project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, max_messages=5, return_immediately=False + project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, max_messages=5, return_immediately=True ) result = operator.get_openlineage_facets_on_complete(operator) 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 cb73498d6a4be..dac40685e2dc2 100644 --- a/providers/google/tests/unit/google/cloud/sensors/test_pubsub.py +++ b/providers/google/tests/unit/google/cloud/sensors/test_pubsub.py @@ -95,7 +95,7 @@ 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=False + project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, max_messages=5, return_immediately=True ) assert generated_dicts == response @@ -145,7 +145,7 @@ def messages_callback( 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 + project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, max_messages=5, return_immediately=True ) messages_callback.assert_called_once() @@ -167,9 +167,9 @@ 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 False + assert exc.value.trigger.return_immediately is True - def test_pubsub_pull_sensor_async_with_return_immediately_true(self): + 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. @@ -180,12 +180,12 @@ def test_pubsub_pull_sensor_async_with_return_immediately_true(self): project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, deferrable=True, - return_immediately=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 True + 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""" 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 f2c0dd4ef7cf7..591aafb5c2db5 100644 --- a/providers/google/tests/unit/google/cloud/triggers/test_pubsub.py +++ b/providers/google/tests/unit/google/cloud/triggers/test_pubsub.py @@ -74,12 +74,12 @@ 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": False, + "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(self, mock_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( @@ -90,7 +90,7 @@ async def test_async_pubsub_pull_trigger_passes_return_immediately(self, mock_pu poke_interval=TEST_POLL_INTERVAL, gcp_conn_id=TEST_GCP_CONN_ID, impersonation_chain=None, - return_immediately=True, + return_immediately=False, ) await trigger.run().asend(None) @@ -99,7 +99,7 @@ async def test_async_pubsub_pull_trigger_passes_return_immediately(self, mock_pu project_id=PROJECT_ID, subscription="subscription", max_messages=MAX_MESSAGES, - return_immediately=True, + return_immediately=False, ) @pytest.mark.asyncio From e9d18e6f6f4473226db3e62ddf718ffa8fc38b2a Mon Sep 17 00:00:00 2001 From: Mike Prieto Date: Mon, 13 Jul 2026 17:29:38 +0000 Subject: [PATCH 4/7] test: Suppress test failures caused by use of the deprecated return_immediately parameter and add better test coverage for the return_immediately parameter --- .../providers/google/cloud/sensors/pubsub.py | 2 +- .../unit/google/cloud/hooks/test_pubsub.py | 25 +++++++++++++++ .../google/cloud/operators/test_pubsub.py | 32 +++++++++++++++++++ .../unit/google/cloud/sensors/test_pubsub.py | 32 +++++++++++++++++++ .../unit/google/cloud/triggers/test_pubsub.py | 16 ++++++++++ 5 files changed, 106 insertions(+), 1 deletion(-) 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 a266f42d76e8a..162393057d20d 100644 --- a/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py @@ -19,10 +19,10 @@ from __future__ import annotations +import warnings from collections.abc import Callable, Sequence from datetime import timedelta from typing import TYPE_CHECKING, Any -import warnings from google.cloud import pubsub_v1 from google.cloud.pubsub_v1.types import ReceivedMessage diff --git a/providers/google/tests/unit/google/cloud/hooks/test_pubsub.py b/providers/google/tests/unit/google/cloud/hooks/test_pubsub.py index eeae2aa763519..6dfc4af3aaf7c 100644 --- a/providers/google/tests/unit/google/cloud/hooks/test_pubsub.py +++ b/providers/google/tests/unit/google/cloud/hooks/test_pubsub.py @@ -28,10 +28,13 @@ from google.cloud.pubsub_v1.types import PublisherOptions, ReceivedMessage from googleapiclient.errors import HttpError +from airflow.exceptions import AirflowProviderDeprecationWarning from airflow.providers.google.cloud.hooks.pubsub import PubSubAsyncHook, PubSubException, PubSubHook from airflow.providers.google.common.consts import CLIENT_INFO from airflow.version import version +pytestmark = pytest.mark.filterwarnings("ignore::airflow.exceptions.AirflowProviderDeprecationWarning") + BASE_STRING = "airflow.providers.google.common.hooks.base_google.{}" PUBSUB_STRING = "airflow.providers.google.cloud.hooks.pubsub.{}" @@ -648,6 +651,17 @@ def test_messages_validation_negative(self, messages, error_message): PubSubHook._validate_messages(messages) assert str(ctx.value) == error_message + @mock.patch("airflow.providers.google.cloud.hooks.pubsub.PubSubHook.subscriber_client") + def test_pull_deprecation_warning(self, mock_subscriber_client): + hook = PubSubHook() + with pytest.warns(AirflowProviderDeprecationWarning, match="return_immediately"): + hook.pull( + project_id=TEST_PROJECT, + subscription=TEST_SUBSCRIPTION, + max_messages=10, + return_immediately=False, + ) + class TestPubSubAsyncHook: @pytest.fixture @@ -696,3 +710,14 @@ async def test_acknowledge(self, mock_subscriber_client, hook): timeout=None, metadata=(), ) + + @pytest.mark.asyncio + @mock.patch("airflow.providers.google.cloud.hooks.pubsub.PubSubAsyncHook._get_subscriber_client") + async def test_pull_deprecation_warning(self, mock_get_subscriber_client, hook): + with pytest.warns(AirflowProviderDeprecationWarning, match="return_immediately"): + await hook.pull( + project_id=TEST_PROJECT, + subscription=TEST_SUBSCRIPTION, + max_messages=10, + return_immediately=False, + ) 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 680aa5d96c694..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, @@ -36,6 +37,8 @@ ) 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" @@ -548,6 +551,26 @@ def test_execute_deferred(self, mock_hook): 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( @@ -650,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 dac40685e2dc2..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( @@ -264,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 591aafb5c2db5..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" @@ -198,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 9ce49fd5b0cd042459b0490a8a2654db24ddfc12 Mon Sep 17 00:00:00 2001 From: Mike Prieto Date: Tue, 21 Jul 2026 13:55:33 +0000 Subject: [PATCH 5/7] Only show the deprecation warning if the return_immediately option is explicitly specified for Pub/Sub modules --- .../providers/google/cloud/hooks/pubsub.py | 30 +++++++++++-------- .../google/cloud/operators/pubsub.py | 22 +++++++------- .../providers/google/cloud/sensors/pubsub.py | 22 +++++++------- .../providers/google/cloud/triggers/pubsub.py | 22 +++++++------- .../unit/google/cloud/hooks/test_pubsub.py | 3 +- 5 files changed, 55 insertions(+), 44 deletions(-) diff --git a/providers/google/src/airflow/providers/google/cloud/hooks/pubsub.py b/providers/google/src/airflow/providers/google/cloud/hooks/pubsub.py index ab1b1cee4c5f7..cbdcb885d7696 100644 --- a/providers/google/src/airflow/providers/google/cloud/hooks/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/hooks/pubsub.py @@ -492,7 +492,7 @@ def pull( subscription: str, max_messages: int, project_id: str = PROVIDE_PROJECT_ID, - return_immediately: bool = False, + return_immediately: bool | None = None, retry: Retry | _MethodDefault = DEFAULT, timeout: float | None = None, metadata: Sequence[tuple[str, str]] = (), @@ -520,11 +520,14 @@ def pull( the base64-encoded message content. See https://cloud.google.com/pubsub/docs/reference/rpc/google.pubsub.v1#google.pubsub.v1.ReceivedMessage """ - warnings.warn( - "The `return_immediately` parameter is deprecated and will be removed in a future release.", - AirflowProviderDeprecationWarning, - stacklevel=2, - ) + if return_immediately is not None: + warnings.warn( + "The `return_immediately` parameter is deprecated and will be removed in a future release.", + AirflowProviderDeprecationWarning, + stacklevel=2, + ) + else: + return_immediately = False subscriber = self.subscriber_client # E501 @@ -690,7 +693,7 @@ async def pull( subscription: str, max_messages: int, project_id: str = PROVIDE_PROJECT_ID, - return_immediately: bool = False, + return_immediately: bool | None = None, retry: AsyncRetry | _MethodDefault = DEFAULT, timeout: float | None = None, metadata: Sequence[tuple[str, str]] = (), @@ -718,11 +721,14 @@ async def pull( the base64-encoded message content. See https://cloud.google.com/pubsub/docs/reference/rpc/google.pubsub.v1#google.pubsub.v1.ReceivedMessage """ - warnings.warn( - "The `return_immediately` parameter is deprecated and will be removed in a future release.", - AirflowProviderDeprecationWarning, - stacklevel=2, - ) + if return_immediately is not None: + warnings.warn( + "The `return_immediately` parameter is deprecated and will be removed in a future release.", + AirflowProviderDeprecationWarning, + stacklevel=2, + ) + else: + return_immediately = False subscriber = await self._get_subscriber_client() subscription_path = f"projects/{project_id}/subscriptions/{subscription}" 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 3cc7c69c19c3f..4d76981980ade 100644 --- a/providers/google/src/airflow/providers/google/cloud/operators/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/operators/pubsub.py @@ -828,7 +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 = True, + return_immediately: bool | None = None, **kwargs, ) -> None: super().__init__(**kwargs) @@ -841,15 +841,17 @@ def __init__( self.impersonation_chain = impersonation_chain self.deferrable = deferrable self.poll_interval = poll_interval - self.return_immediately = return_immediately - - warnings.warn( - "The `return_immediately` parameter is deprecated and will be removed in a future release. " - "Its default value will be changed to `False` in the next major release. " - "Planned removal date: August 01, 2026.", - AirflowProviderDeprecationWarning, - stacklevel=2, - ) + if return_immediately is not None: + warnings.warn( + "The `return_immediately` parameter is deprecated and will be removed in a future release. " + "Its default value will be changed to `False` in the next major release. " + "Planned removal date: August 01, 2026.", + AirflowProviderDeprecationWarning, + stacklevel=2, + ) + self.return_immediately = return_immediately + else: + self.return_immediately = True 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 162393057d20d..93be5ce5978e1 100644 --- a/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py @@ -115,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, @@ -129,21 +129,23 @@ 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 - - warnings.warn( - "The `return_immediately` parameter is deprecated and will be removed in a future release. " - "Its default value will be changed to `False` in the next major release. " - "Planned removal date: August 01, 2026.", - AirflowProviderDeprecationWarning, - stacklevel=2, - ) + if return_immediately is not None: + warnings.warn( + "The `return_immediately` parameter is deprecated and will be removed in a future release. " + "Its default value will be changed to `False` in the next major release. " + "Planned removal date: August 01, 2026.", + 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 7e87f0c9039df..aa1a6c3f4952c 100644 --- a/providers/google/src/airflow/providers/google/cloud/triggers/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/triggers/pubsub.py @@ -75,7 +75,7 @@ def __init__( gcp_conn_id: str, poke_interval: float = 10.0, impersonation_chain: str | Sequence[str] | None = None, - return_immediately: bool = True, + return_immediately: bool | None = None, ): super().__init__() self.project_id = project_id @@ -85,15 +85,17 @@ def __init__( self.poke_interval = poke_interval self.gcp_conn_id = gcp_conn_id self.impersonation_chain = impersonation_chain - self.return_immediately = return_immediately - - warnings.warn( - "The `return_immediately` parameter is deprecated and will be removed in a future release. " - "Its default value will be changed to `False` in the next major release. " - "Planned removal date: August 01, 2026.", - AirflowProviderDeprecationWarning, - stacklevel=2, - ) + if return_immediately is not None: + warnings.warn( + "The `return_immediately` parameter is deprecated and will be removed in a future release. " + "Its default value will be changed to `False` in the next major release. " + "Planned removal date: August 01, 2026.", + 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.""" diff --git a/providers/google/tests/unit/google/cloud/hooks/test_pubsub.py b/providers/google/tests/unit/google/cloud/hooks/test_pubsub.py index 6dfc4af3aaf7c..84017146f08bb 100644 --- a/providers/google/tests/unit/google/cloud/hooks/test_pubsub.py +++ b/providers/google/tests/unit/google/cloud/hooks/test_pubsub.py @@ -653,9 +653,8 @@ def test_messages_validation_negative(self, messages, error_message): @mock.patch("airflow.providers.google.cloud.hooks.pubsub.PubSubHook.subscriber_client") def test_pull_deprecation_warning(self, mock_subscriber_client): - hook = PubSubHook() with pytest.warns(AirflowProviderDeprecationWarning, match="return_immediately"): - hook.pull( + self.pubsub_hook.pull( project_id=TEST_PROJECT, subscription=TEST_SUBSCRIPTION, max_messages=10, From 478e40018c77bb77d01061103b41e2cd861ca13a Mon Sep 17 00:00:00 2001 From: Mike Prieto Date: Thu, 23 Jul 2026 20:48:41 +0000 Subject: [PATCH 6/7] Update deprecation warning to only include plan to change default value of return_immediately to False and undo all changes to PubSubHook as it already defaults return_immediately to False --- .../providers/google/cloud/hooks/pubsub.py | 22 ++--------------- .../google/cloud/operators/pubsub.py | 4 +--- .../providers/google/cloud/sensors/pubsub.py | 4 +--- .../providers/google/cloud/triggers/pubsub.py | 4 +--- .../unit/google/cloud/hooks/test_pubsub.py | 24 ------------------- 5 files changed, 5 insertions(+), 53 deletions(-) diff --git a/providers/google/src/airflow/providers/google/cloud/hooks/pubsub.py b/providers/google/src/airflow/providers/google/cloud/hooks/pubsub.py index cbdcb885d7696..a983fe3e4b7c7 100644 --- a/providers/google/src/airflow/providers/google/cloud/hooks/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/hooks/pubsub.py @@ -492,7 +492,7 @@ def pull( subscription: str, max_messages: int, project_id: str = PROVIDE_PROJECT_ID, - return_immediately: bool | None = None, + return_immediately: bool = False, retry: Retry | _MethodDefault = DEFAULT, timeout: float | None = None, metadata: Sequence[tuple[str, str]] = (), @@ -520,15 +520,6 @@ def pull( the base64-encoded message content. See https://cloud.google.com/pubsub/docs/reference/rpc/google.pubsub.v1#google.pubsub.v1.ReceivedMessage """ - if return_immediately is not None: - warnings.warn( - "The `return_immediately` parameter is deprecated and will be removed in a future release.", - AirflowProviderDeprecationWarning, - stacklevel=2, - ) - else: - return_immediately = False - subscriber = self.subscriber_client # E501 subscription_path = f"projects/{project_id}/subscriptions/{subscription}" @@ -693,7 +684,7 @@ async def pull( subscription: str, max_messages: int, project_id: str = PROVIDE_PROJECT_ID, - return_immediately: bool | None = None, + return_immediately: bool = False, retry: AsyncRetry | _MethodDefault = DEFAULT, timeout: float | None = None, metadata: Sequence[tuple[str, str]] = (), @@ -721,15 +712,6 @@ async def pull( the base64-encoded message content. See https://cloud.google.com/pubsub/docs/reference/rpc/google.pubsub.v1#google.pubsub.v1.ReceivedMessage """ - if return_immediately is not None: - warnings.warn( - "The `return_immediately` parameter is deprecated and will be removed in a future release.", - AirflowProviderDeprecationWarning, - stacklevel=2, - ) - else: - return_immediately = False - subscriber = await self._get_subscriber_client() subscription_path = f"projects/{project_id}/subscriptions/{subscription}" self.log.info("Pulling max %d messages from subscription (path) %s", max_messages, subscription_path) 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 4d76981980ade..e4352fa063214 100644 --- a/providers/google/src/airflow/providers/google/cloud/operators/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/operators/pubsub.py @@ -843,9 +843,7 @@ def __init__( self.poll_interval = poll_interval if return_immediately is not None: warnings.warn( - "The `return_immediately` parameter is deprecated and will be removed in a future release. " - "Its default value will be changed to `False` in the next major release. " - "Planned removal date: August 01, 2026.", + "The default value of `return_immediately` will be changed to `False` in a future release.", AirflowProviderDeprecationWarning, stacklevel=2, ) 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 93be5ce5978e1..25a4c618fc113 100644 --- a/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py @@ -137,9 +137,7 @@ def __init__( self._return_value = None if return_immediately is not None: warnings.warn( - "The `return_immediately` parameter is deprecated and will be removed in a future release. " - "Its default value will be changed to `False` in the next major release. " - "Planned removal date: August 01, 2026.", + "The default value of `return_immediately` will be changed to `False` in a future release.", AirflowProviderDeprecationWarning, stacklevel=2, ) 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 aa1a6c3f4952c..fce2d9b8c04ef 100644 --- a/providers/google/src/airflow/providers/google/cloud/triggers/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/triggers/pubsub.py @@ -87,9 +87,7 @@ def __init__( self.impersonation_chain = impersonation_chain if return_immediately is not None: warnings.warn( - "The `return_immediately` parameter is deprecated and will be removed in a future release. " - "Its default value will be changed to `False` in the next major release. " - "Planned removal date: August 01, 2026.", + "The default value of `return_immediately` will be changed to `False` in a future release.", AirflowProviderDeprecationWarning, stacklevel=2, ) diff --git a/providers/google/tests/unit/google/cloud/hooks/test_pubsub.py b/providers/google/tests/unit/google/cloud/hooks/test_pubsub.py index 84017146f08bb..eeae2aa763519 100644 --- a/providers/google/tests/unit/google/cloud/hooks/test_pubsub.py +++ b/providers/google/tests/unit/google/cloud/hooks/test_pubsub.py @@ -28,13 +28,10 @@ from google.cloud.pubsub_v1.types import PublisherOptions, ReceivedMessage from googleapiclient.errors import HttpError -from airflow.exceptions import AirflowProviderDeprecationWarning from airflow.providers.google.cloud.hooks.pubsub import PubSubAsyncHook, PubSubException, PubSubHook from airflow.providers.google.common.consts import CLIENT_INFO from airflow.version import version -pytestmark = pytest.mark.filterwarnings("ignore::airflow.exceptions.AirflowProviderDeprecationWarning") - BASE_STRING = "airflow.providers.google.common.hooks.base_google.{}" PUBSUB_STRING = "airflow.providers.google.cloud.hooks.pubsub.{}" @@ -651,16 +648,6 @@ def test_messages_validation_negative(self, messages, error_message): PubSubHook._validate_messages(messages) assert str(ctx.value) == error_message - @mock.patch("airflow.providers.google.cloud.hooks.pubsub.PubSubHook.subscriber_client") - def test_pull_deprecation_warning(self, mock_subscriber_client): - with pytest.warns(AirflowProviderDeprecationWarning, match="return_immediately"): - self.pubsub_hook.pull( - project_id=TEST_PROJECT, - subscription=TEST_SUBSCRIPTION, - max_messages=10, - return_immediately=False, - ) - class TestPubSubAsyncHook: @pytest.fixture @@ -709,14 +696,3 @@ async def test_acknowledge(self, mock_subscriber_client, hook): timeout=None, metadata=(), ) - - @pytest.mark.asyncio - @mock.patch("airflow.providers.google.cloud.hooks.pubsub.PubSubAsyncHook._get_subscriber_client") - async def test_pull_deprecation_warning(self, mock_get_subscriber_client, hook): - with pytest.warns(AirflowProviderDeprecationWarning, match="return_immediately"): - await hook.pull( - project_id=TEST_PROJECT, - subscription=TEST_SUBSCRIPTION, - max_messages=10, - return_immediately=False, - ) From f2d9b116e558f91baf5fc57b94426b23eb227b27 Mon Sep 17 00:00:00 2001 From: Mike Prieto Date: Thu, 23 Jul 2026 20:54:14 +0000 Subject: [PATCH 7/7] Update wording for return_immediately default value deprecation warning --- .../src/airflow/providers/google/cloud/operators/pubsub.py | 2 +- .../google/src/airflow/providers/google/cloud/sensors/pubsub.py | 2 +- .../src/airflow/providers/google/cloud/triggers/pubsub.py | 2 +- 3 files changed, 3 insertions(+), 3 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 e4352fa063214..207ca834004b0 100644 --- a/providers/google/src/airflow/providers/google/cloud/operators/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/operators/pubsub.py @@ -843,7 +843,7 @@ def __init__( 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 release.", + "The default value of `return_immediately` will be changed to `False` in a future major release.", AirflowProviderDeprecationWarning, stacklevel=2, ) 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 25a4c618fc113..7c9ceccd2774a 100644 --- a/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/sensors/pubsub.py @@ -137,7 +137,7 @@ def __init__( 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 release.", + "The default value of `return_immediately` will be changed to `False` in a future major release.", AirflowProviderDeprecationWarning, stacklevel=2, ) 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 fce2d9b8c04ef..4c8aff5d27ac4 100644 --- a/providers/google/src/airflow/providers/google/cloud/triggers/pubsub.py +++ b/providers/google/src/airflow/providers/google/cloud/triggers/pubsub.py @@ -87,7 +87,7 @@ def __init__( 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 release.", + "The default value of `return_immediately` will be changed to `False` in a future major release.", AirflowProviderDeprecationWarning, stacklevel=2, )