This is an automated email from the ASF dual-hosted git repository.
eladkal pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 57183056ca7 Fix AsyncToSync error in Databricks deferrable operators
on Airflow 3.0 (#74406)
57183056ca7 is described below
commit 57183056ca7ba5e99163ff6d686c4de4f5911b8c
Author: Noritaka Sekiyama <[email protected]>
AuthorDate: Thu Oct 8 02:59:43 2026 +0900
Fix AsyncToSync error in Databricks deferrable operators on Airflow 3.0
(#74406)
* Fix AsyncToSync error in Databricks deferrable operators on Airflow 3.0
On Airflow 3.0 the triggerer cannot run the synchronous connection lookup
on its event loop, so every deferrable Databricks operator failed as soon
as its trigger made the first API call. The connection is now fetched
asynchronously once per hook, which also keeps the existing sync helpers
usable from the async code path.
The hook argument of get_async_connection first shipped in
common-compat 1.17.0, hence the raised floor.
closes: #71525
Co-authored-by: Saiteja Bandaru <[email protected]>
Co-authored-by: Isaac <[email protected]>
* Leave the common-compat floor bump to the release manager
Only release managers may raise provider dependency floors; the
"use next version" marker asks them to pin the next common-compat
release, which already includes the hook argument of
get_async_connection.
Co-authored-by: Isaac <[email protected]>
---------
Co-authored-by: Saiteja Bandaru <[email protected]>
Co-authored-by: Isaac <[email protected]>
---
providers/databricks/pyproject.toml | 2 +-
.../providers/databricks/hooks/databricks_base.py | 8 ++++++++
.../tests/unit/databricks/hooks/test_databricks.py | 17 +++++++++++++++++
3 files changed, 26 insertions(+), 1 deletion(-)
diff --git a/providers/databricks/pyproject.toml
b/providers/databricks/pyproject.toml
index 5b29e355056..2bfe0183794 100644
--- a/providers/databricks/pyproject.toml
+++ b/providers/databricks/pyproject.toml
@@ -59,7 +59,7 @@ requires-python = ">=3.11"
# After you modify the dependencies, and rebuild your Breeze CI image with
``breeze ci-image build``
dependencies = [
"apache-airflow>=2.11.0",
- "apache-airflow-providers-common-compat>=1.13.0",
+ "apache-airflow-providers-common-compat>=1.13.0", # use next version
"apache-airflow-providers-common-sql>=1.32.0",
"requests>=2.32.0,<3",
"databricks-sql-connector>=4.4.0",
diff --git
a/providers/databricks/src/airflow/providers/databricks/hooks/databricks_base.py
b/providers/databricks/src/airflow/providers/databricks/hooks/databricks_base.py
index fe028f881e2..d5dcedaa451 100644
---
a/providers/databricks/src/airflow/providers/databricks/hooks/databricks_base.py
+++
b/providers/databricks/src/airflow/providers/databricks/hooks/databricks_base.py
@@ -51,6 +51,7 @@ from tenacity import (
)
from airflow import __version__
+from airflow.providers.common.compat.connection import get_async_connection
from airflow.providers.common.compat.module_loading import import_string
from airflow.providers.common.compat.sdk import AirflowException,
AirflowOptionalProviderFeatureException
from airflow.providers.databricks.exceptions import DatabricksApiError
@@ -197,6 +198,12 @@ class BaseDatabricksHook(BaseHook):
def get_conn(self) -> Connection:
return self.databricks_conn
+ async def _a_cache_databricks_conn(self) -> None:
+ # The sync ``get_connection`` cannot run on the triggerer's event loop
on Airflow 3.0, so fill
+ # the ``databricks_conn`` cache asynchronously and let the sync
helpers read it from there.
+ if "databricks_conn" not in self.__dict__:
+ self.__dict__["databricks_conn"] = await
get_async_connection(self.databricks_conn_id, hook=self)
+
@cached_property
def user_agent_header(self) -> dict[str, str]:
return {"user-agent": self.user_agent_value}
@@ -1393,6 +1400,7 @@ class BaseDatabricksHook(BaseHook):
:return: If the api call returns a OK status code,
this function returns the response in JSON. Otherwise, throw an
AirflowException.
"""
+ await self._a_cache_databricks_conn()
method, endpoint = endpoint_info
full_endpoint = f"api/{endpoint}"
diff --git
a/providers/databricks/tests/unit/databricks/hooks/test_databricks.py
b/providers/databricks/tests/unit/databricks/hooks/test_databricks.py
index 38d0bd1da10..06c588b2df3 100644
--- a/providers/databricks/tests/unit/databricks/hooks/test_databricks.py
+++ b/providers/databricks/tests/unit/databricks/hooks/test_databricks.py
@@ -1569,6 +1569,23 @@ class
TestDatabricksHookConnSettings(TestDatabricksHookToken):
mock_get.assert_called_once()
assert mock_get.call_args.args ==
(f"http://{HOST}:7908/api/2.1/foo/bar",)
+ @pytest.mark.asyncio
+ @mock.patch.object(
+ DatabricksHook,
+ "get_connection",
+ autospec=True,
+ side_effect=RuntimeError("You cannot use AsyncToSync in the same
thread as an async event loop"),
+ )
+
@mock.patch("airflow.providers.databricks.hooks.databricks_base.aiohttp.ClientSession.get")
+ async def test_async_do_api_call_fetches_connection_asynchronously(self,
mock_get, mock_get_connection):
+ mock_get.return_value.__aenter__.return_value.json =
AsyncMock(return_value={"bar": "baz"})
+ async with self.hook:
+ run_page_url = await self.hook._a_do_api_call(("GET",
"2.1/foo/bar"))
+
+ assert run_page_url == {"bar": "baz"}
+ assert mock_get.call_args.args ==
(f"http://{HOST}:7908/api/2.1/foo/bar",)
+ mock_get_connection.assert_not_called()
+
@pytest.mark.asyncio
@mock.patch("airflow.providers.databricks.hooks.databricks_base.aiohttp.ClientSession.get")
async def
test_async_do_api_call_only_existing_response_properties_are_read(self,
mock_get):