Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,8 @@
from collections.abc import AsyncIterator
from typing import Any

from airflow.providers.common.compat.sdk import Connection


class SQLExecuteQueryTrigger(BaseTrigger):
"""
Expand Down Expand Up @@ -61,13 +63,22 @@ def serialize(self) -> tuple[str, dict[str, Any]]:
},
)

def get_hook(self) -> DbApiHook:
@classmethod
async def get_async_connection(cls, conn_id: str) -> Connection:
if hasattr(BaseHook, "aget_connection"):
return await BaseHook.aget_connection(conn_id=conn_id)

from asgiref.sync import sync_to_async

return await sync_to_async(BaseHook.get_connection)(conn_id=conn_id)

async def get_hook(self) -> DbApiHook:
"""
Return DbApiHook.

:return: DbApiHook for this connection
"""
connection = BaseHook.get_connection(self.conn_id)
connection = await self.get_async_connection(conn_id=self.conn_id)
hook = connection.get_hook(hook_params=self.hook_params)
if not isinstance(hook, DbApiHook):
raise AirflowException(
Expand All @@ -81,7 +92,7 @@ def get_hook(self) -> DbApiHook:

async def run(self) -> AsyncIterator[TriggerEvent]:
try:
hook = self.get_hook()
hook = await self.get_hook()

self.log.info("Extracting data from %s", self.conn_id)
self.log.info("Executing: \n %s", self.sql)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,13 @@ def get_base_airflow_version_tuple() -> tuple[int, int, int]:
AIRFLOW_V_3_0_PLUS = get_base_airflow_version_tuple() >= (3, 0, 0)
AIRFLOW_V_3_1_PLUS: bool = get_base_airflow_version_tuple() >= (3, 1, 0)

if AIRFLOW_V_3_1_PLUS:
from airflow.sdk import BaseHook
else:
from airflow.hooks.base import BaseHook # type: ignore[attr-defined,no-redef]

__all__ = [
"AIRFLOW_V_3_0_PLUS",
"AIRFLOW_V_3_1_PLUS",
"BaseHook",
]
Original file line number Diff line number Diff line change
Expand Up @@ -24,12 +24,14 @@
from unittest.mock import MagicMock

import pytest
from asgiref.sync import sync_to_async
from more_itertools import flatten

from airflow.exceptions import AirflowProviderDeprecationWarning
from airflow.models.connection import Connection
from airflow.models.dag import DAG
from airflow.providers.common.sql.hooks.sql import DbApiHook
from airflow.providers.common.sql.version_compat import BaseHook
from airflow.providers.mysql.hooks.mysql import MySqlHook
from airflow.providers.postgres.hooks.postgres import PostgresHook
from airflow.utils import timezone
Expand All @@ -38,15 +40,6 @@
from tests_common.test_utils.operators.run_deferrable import execute_operator, mock_context
from tests_common.test_utils.providers import get_provider_min_airflow_version

try:
import importlib.util

if not importlib.util.find_spec("airflow.sdk.bases.hook"):
raise ImportError

BASEHOOK_PATCH_PATH = "airflow.sdk.bases.hook.BaseHook"
except ImportError:
BASEHOOK_PATCH_PATH = "airflow.hooks.base.BaseHook"
pytestmark = pytest.mark.db_test

DEFAULT_DATE = timezone.datetime(2015, 1, 1)
Expand Down Expand Up @@ -267,19 +260,21 @@ def test_templated_fields(self):
assert operator.paginated_sql_statement_clause == "{} OFFSET {} ROWS FETCH NEXT {} ROWS ONLY;"

def test_non_paginated_read(self):
with mock.patch(f"{BASEHOOK_PATCH_PATH}.get_connection", side_effect=self.get_connection):
with mock.patch(f"{BASEHOOK_PATCH_PATH}.get_hook", side_effect=self.get_hook):
operator = GenericTransfer(
task_id="transfer_table",
source_conn_id="my_source_conn_id",
destination_conn_id="my_destination_conn_id",
sql="SELECT * FROM HR.EMPLOYEES",
destination_table="NEW_HR.EMPLOYEES",
insert_args=INSERT_ARGS,
execution_timeout=timedelta(hours=1),
)
with (
mock.patch.object(BaseHook, "get_connection", side_effect=self.get_connection),
mock.patch.object(BaseHook, "get_hook", side_effect=self.get_hook),
):
operator = GenericTransfer(
task_id="transfer_table",
source_conn_id="my_source_conn_id",
destination_conn_id="my_destination_conn_id",
sql="SELECT * FROM HR.EMPLOYEES",
destination_table="NEW_HR.EMPLOYEES",
insert_args=INSERT_ARGS,
execution_timeout=timedelta(hours=1),
)

operator.execute(context=mock_context(task=operator))
operator.execute(context=mock_context(task=operator))

assert self.mocked_source_hook.get_records.call_count == 1
assert self.mocked_source_hook.get_records.call_args_list[0].args[0] == "SELECT * FROM HR.EMPLOYEES"
Expand Down Expand Up @@ -327,26 +322,29 @@ def test_paginated_read(self):
https://medium.com/apache-airflow/transfering-data-from-sap-hana-to-mssql-using-the-airflow-generictransfer-d29f147a9f1f
"""

with mock.patch(f"{BASEHOOK_PATCH_PATH}.get_connection", side_effect=self.get_connection):
with mock.patch(f"{BASEHOOK_PATCH_PATH}.get_hook", side_effect=self.get_hook):
operator = GenericTransfer(
task_id="transfer_table",
source_conn_id="my_source_conn_id",
destination_conn_id="my_destination_conn_id",
sql="SELECT * FROM HR.EMPLOYEES",
destination_table="NEW_HR.EMPLOYEES",
page_size=1000, # Fetch data in chunks of 1000 rows for pagination
insert_args=INSERT_ARGS,
execution_timeout=timedelta(hours=1),
)
with (
mock.patch.object(BaseHook, "get_connection", side_effect=self.get_connection),
mock.patch.object(BaseHook, "aget_connection", side_effect=sync_to_async(self.get_connection)),
mock.patch.object(BaseHook, "get_hook", side_effect=self.get_hook),
):
operator = GenericTransfer(
task_id="transfer_table",
source_conn_id="my_source_conn_id",
destination_conn_id="my_destination_conn_id",
sql="SELECT * FROM HR.EMPLOYEES",
destination_table="NEW_HR.EMPLOYEES",
page_size=1000, # Fetch data in chunks of 1000 rows for pagination
insert_args=INSERT_ARGS,
execution_timeout=timedelta(hours=1),
)

results, events = execute_operator(operator)
results, events = execute_operator(operator)

assert not results
assert len(events) == 3
assert events[0].payload["results"] == [[1, 2], [11, 12], [3, 4], [13, 14]]
assert events[1].payload["results"] == [[3, 4], [13, 14]]
assert not events[2].payload["results"]
assert not results
assert len(events) == 3
assert events[0].payload["results"] == [[1, 2], [11, 12], [3, 4], [13, 14]]
assert events[1].payload["results"] == [[3, 4], [13, 14]]
assert not events[2].payload["results"]

assert self.mocked_source_hook.get_records.call_count == 3
assert (
Expand Down
55 changes: 34 additions & 21 deletions providers/common/sql/tests/unit/common/sql/triggers/test_sql.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,38 +17,51 @@
from __future__ import annotations

from unittest import mock
from unittest.mock import AsyncMock

import pytest

from airflow.models.connection import Connection
from airflow.providers.common.sql.hooks.sql import DbApiHook
from airflow.providers.common.sql.triggers.sql import SQLExecuteQueryTrigger
from airflow.providers.common.sql.version_compat import BaseHook
from airflow.triggers.base import TriggerEvent

try:
import importlib.util

if not importlib.util.find_spec("airflow.sdk.bases.hook"):
raise ImportError

BASEHOOK_PATCH_PATH = "airflow.sdk.bases.hook.BaseHook"
except ImportError:
BASEHOOK_PATCH_PATH = "airflow.hooks.base.BaseHook"
from tests_common.test_utils.operators.run_deferrable import run_trigger


class TestSQLExecuteQueryTrigger:
@mock.patch(f"{BASEHOOK_PATCH_PATH}.get_connection")
def test_run(self, mock_get_connection):
data = [(1, "Alice"), (2, "Bob")]
mock_connection = mock.MagicMock(spec=Connection)
@classmethod
def get_connection(cls, conn_id: str | None = None, records: list | None = None):
mock_connection = mock.MagicMock(spec=Connection, conn_id=conn_id)
mock_hook = mock.MagicMock(spec=DbApiHook)
mock_hook.get_records.side_effect = lambda sql: data
mock_get_connection.return_value = mock_connection
if records:
mock_hook.get_records.side_effect = lambda sql: records
mock_connection.get_hook.side_effect = lambda hook_params: mock_hook
return mock_connection

def test_run(self):
data = [(1, "Alice"), (2, "Bob")]

with (
mock.patch.object(BaseHook, "get_connection", side_effect=lambda conn_id: self.get_connection(conn_id, data)),
mock.patch.object(BaseHook, "aget_connection", new_callable=lambda: AsyncMock(side_effect=lambda *args, **kwargs: self.get_connection(kwargs.get("conn_id") or args[0], data))),
):
trigger = SQLExecuteQueryTrigger(sql="SELECT * FROM users;", conn_id="test_conn_id")
actual = run_trigger(trigger)

assert len(actual) == 1
assert isinstance(actual[0], TriggerEvent)
assert actual[0].payload["status"] == "success"
assert actual[0].payload["results"] == data

trigger = SQLExecuteQueryTrigger(sql="SELECT * FROM users;", conn_id="test_conn_id")
actual = run_trigger(trigger)
@pytest.mark.asyncio
async def test_get_hook(self):
with (
mock.patch.object(BaseHook, "get_connection", side_effect=lambda conn_id: self.get_connection(conn_id)),
mock.patch.object(BaseHook, "aget_connection", new_callable=lambda: AsyncMock(side_effect=lambda *args, **kwargs: self.get_connection(kwargs.get("conn_id") or args[0]))),
):
trigger = SQLExecuteQueryTrigger(sql="SELECT * FROM users;", conn_id="test_conn_id")
actual = await trigger.get_hook()

assert len(actual) == 1
assert isinstance(actual[0], TriggerEvent)
assert actual[0].payload["status"] == "success"
assert actual[0].payload["results"] == data
assert isinstance(actual, DbApiHook)
Loading