This is an automated email from the ASF dual-hosted git repository.
shahar1 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 c1e98dc0c96 Add dedicated InfluxDB 3 sensor (#73522)
c1e98dc0c96 is described below
commit c1e98dc0c96eb9177f6086a4aa551d0200b65377
Author: Subhramit Basu <[email protected]>
AuthorDate: Fri Sep 25 12:41:11 2026 +0530
Add dedicated InfluxDB 3 sensor (#73522)
---
providers/influxdb/docs/index.rst | 1 +
providers/influxdb/docs/sensors/index.rst | 40 +++++
providers/influxdb/provider.yaml | 5 +
.../providers/influxdb/get_provider_info.py | 6 +
.../influxdb/{utils.py => sensors/__init__.py} | 14 --
.../providers/influxdb/sensors/influxdb3.py | 108 +++++++++++++
.../providers/influxdb/triggers/influxdb3.py | 52 ++++++-
.../src/airflow/providers/influxdb/utils.py | 16 ++
.../tests/system/influxdb/example_influxdb3.py | 14 +-
.../unit/influxdb/sensors/__init__.py} | 14 --
.../tests/unit/influxdb/sensors/test_influxdb3.py | 170 +++++++++++++++++++++
.../influxdb/tests/unit/influxdb/test_utils.py | 29 +++-
.../tests/unit/influxdb/triggers/test_influxdb3.py | 123 ++++++++++++++-
13 files changed, 560 insertions(+), 32 deletions(-)
diff --git a/providers/influxdb/docs/index.rst
b/providers/influxdb/docs/index.rst
index dfc9b71ad07..634405ee527 100644
--- a/providers/influxdb/docs/index.rst
+++ b/providers/influxdb/docs/index.rst
@@ -37,6 +37,7 @@
Connection types (InfluxDB 2.x) <connections/influxdb>
Connection types (InfluxDB 3.x) <connections/influxdb3>
Operators <operators/index>
+ Sensors <sensors/index>
.. toctree::
:hidden:
diff --git a/providers/influxdb/docs/sensors/index.rst
b/providers/influxdb/docs/sensors/index.rst
new file mode 100644
index 00000000000..dc81fa62663
--- /dev/null
+++ b/providers/influxdb/docs/sensors/index.rst
@@ -0,0 +1,40 @@
+.. Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+.. http://www.apache.org/licenses/LICENSE-2.0
+
+.. Unless required by applicable law or agreed to in writing,
+ software distributed under the License is distributed on an
+ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ KIND, either express or implied. See the License for the
+ specific language governing permissions and limitations
+ under the License.
+
+.. _howto/sensor:InfluxDB3Sensor:
+
+InfluxDB3Sensor
+===============
+
+Use :class:`~airflow.providers.influxdb.sensors.influxdb3.InfluxDB3Sensor` to
wait until an
+InfluxDB 3.x SQL query returns a truthy first cell. Prefer an efficient
existence query that returns
+one value and limits the result to one row.
+
+.. exampleinclude:: /../../influxdb/tests/system/influxdb/example_influxdb3.py
+ :language: python
+ :start-after: [START howto_sensor_influxdb3]
+ :end-before: [END howto_sensor_influxdb3]
+
+An empty result, a missing value, numeric or string zero, and an empty string
are treated as false.
+Set ``fail_on_empty=True`` to fail immediately when the query returns no rows.
+
+Deferrable mode
+^^^^^^^^^^^^^^^
+
+Set ``deferrable=True`` to release the worker slot between queries. The
+:class:`~airflow.providers.influxdb.triggers.influxdb3.InfluxDB3SensorTrigger`
repeats the query
+at the configured ``poke_interval`` until the condition is met or Airflow
reaches the sensor timeout.
diff --git a/providers/influxdb/provider.yaml b/providers/influxdb/provider.yaml
index b637a1d2f9a..ddc3d971df5 100644
--- a/providers/influxdb/provider.yaml
+++ b/providers/influxdb/provider.yaml
@@ -92,6 +92,11 @@ operators:
python-modules:
- airflow.providers.influxdb.operators.influxdb3
+sensors:
+ - integration-name: InfluxDB 3
+ python-modules:
+ - airflow.providers.influxdb.sensors.influxdb3
+
triggers:
- integration-name: InfluxDB 3
python-modules:
diff --git
a/providers/influxdb/src/airflow/providers/influxdb/get_provider_info.py
b/providers/influxdb/src/airflow/providers/influxdb/get_provider_info.py
index 0b2011da40f..2fbde52a3f9 100644
--- a/providers/influxdb/src/airflow/providers/influxdb/get_provider_info.py
+++ b/providers/influxdb/src/airflow/providers/influxdb/get_provider_info.py
@@ -56,6 +56,12 @@ def get_provider_info():
"python-modules":
["airflow.providers.influxdb.operators.influxdb3"],
},
],
+ "sensors": [
+ {
+ "integration-name": "InfluxDB 3",
+ "python-modules":
["airflow.providers.influxdb.sensors.influxdb3"],
+ }
+ ],
"triggers": [
{
"integration-name": "InfluxDB 3",
diff --git a/providers/influxdb/src/airflow/providers/influxdb/utils.py
b/providers/influxdb/src/airflow/providers/influxdb/sensors/__init__.py
similarity index 63%
copy from providers/influxdb/src/airflow/providers/influxdb/utils.py
copy to providers/influxdb/src/airflow/providers/influxdb/sensors/__init__.py
index 79abb8d6ee2..13a83393a91 100644
--- a/providers/influxdb/src/airflow/providers/influxdb/utils.py
+++ b/providers/influxdb/src/airflow/providers/influxdb/sensors/__init__.py
@@ -14,17 +14,3 @@
# KIND, either express or implied. See the License for the
# specific language governing permissions and limitations
# under the License.
-"""Internal helpers shared across the InfluxDB provider."""
-
-from __future__ import annotations
-
-import json
-from typing import TYPE_CHECKING, Any
-
-if TYPE_CHECKING:
- import pandas as pd
-
-
-def _convert_dataframe_to_records(dataframe: pd.DataFrame) -> list[dict[str,
Any]]:
- """Convert a query result DataFrame into a JSON-serializable list of
dictionaries."""
- return json.loads(dataframe.to_json(orient="records", date_format="iso"))
diff --git
a/providers/influxdb/src/airflow/providers/influxdb/sensors/influxdb3.py
b/providers/influxdb/src/airflow/providers/influxdb/sensors/influxdb3.py
new file mode 100644
index 00000000000..88067a0bb48
--- /dev/null
+++ b/providers/influxdb/src/airflow/providers/influxdb/sensors/influxdb3.py
@@ -0,0 +1,108 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""Sensor waiting for a SQL query to return a truthy first cell in InfluxDB
3.x."""
+
+from __future__ import annotations
+
+from collections.abc import Sequence
+from datetime import timedelta
+from typing import TYPE_CHECKING, Any
+
+from airflow.providers.common.compat.sdk import AirflowFailException,
BaseSensorOperator, conf
+from airflow.providers.influxdb.hooks.influxdb3 import InfluxDB3Hook
+from airflow.providers.influxdb.triggers.influxdb3 import
InfluxDB3SensorTrigger
+from airflow.providers.influxdb.utils import _first_cell_is_truthy
+
+if TYPE_CHECKING:
+ from airflow.sdk.definitions.context import Context
+
+
+class InfluxDB3Sensor(BaseSensorOperator):
+ """
+ Wait until an InfluxDB 3.x SQL query returns a truthy first cell.
+
+ .. seealso::
+ For more information on how to use this sensor, take a look at the
guide:
+ :ref:`howto/sensor:InfluxDB3Sensor`
+
+ :param sql: The SQL query to poll.
+ :param influxdb3_conn_id: Reference to :ref:`InfluxDB 3 connection id
<howto/connection:influxdb3>`.
+ Defaults to ``influxdb3_default``.
+ :param fail_on_empty: Fail instead of waiting when the query returns no
rows. Defaults to ``False``.
+ :param deferrable: Run polling in the triggerer. Defaults to the
+ ``operators.default_deferrable`` configuration (``False`` if unset).
+ """
+
+ template_fields: Sequence[str] = ("sql", "influxdb3_conn_id")
+ template_ext: Sequence[str] = (".sql",)
+
+ def __init__(
+ self,
+ *,
+ sql: str,
+ influxdb3_conn_id: str = "influxdb3_default",
+ fail_on_empty: bool = False,
+ deferrable: bool = conf.getboolean("operators", "default_deferrable",
fallback=False),
+ **kwargs,
+ ) -> None:
+ super().__init__(**kwargs)
+ self.sql = sql
+ self.influxdb3_conn_id = influxdb3_conn_id
+ self.fail_on_empty = fail_on_empty
+ self.deferrable = deferrable
+
+ def poke(self, context: Context) -> bool:
+ """Return whether the query result meets the sensor condition."""
+ self.log.info("Poking with SQL query: %s", self.sql)
+ dataframe =
InfluxDB3Hook(conn_id=self.influxdb3_conn_id).query(self.sql)
+ if dataframe.empty and self.fail_on_empty:
+ raise AirflowFailException("No rows returned, raising as per
fail_on_empty flag")
+ return _first_cell_is_truthy(dataframe)
+
+ def execute(self, context: Context) -> None:
+ if not self.deferrable:
+ super().execute(context)
+ return
+
+ if self.poke(context):
+ return
+
+ self.defer(
+ timeout=timedelta(seconds=self.timeout),
+ trigger=InfluxDB3SensorTrigger(
+ sql=self.sql,
+ influxdb3_conn_id=self.influxdb3_conn_id,
+ poll_interval=self.poke_interval,
+ fail_on_empty=self.fail_on_empty,
+ ),
+ method_name="execute_complete",
+ )
+
+ def execute_complete(self, context: Context, event: dict[str, Any] | None
= None) -> None:
+ """Complete after the trigger reports that the condition was met."""
+ if event is None:
+ raise RuntimeError("InfluxDB 3 sensor did not return an event")
+
+ status = event.get("status")
+ if status == "fail":
+ raise AirflowFailException(event.get("message", "InfluxDB 3 sensor
failed"))
+ if status == "error":
+ raise RuntimeError(event.get("message", "InfluxDB 3 sensor
failed"))
+ if status != "success":
+ raise RuntimeError(f"InfluxDB 3 sensor returned unexpected status:
{status!r}")
+
+ self.log.info("InfluxDB 3 sensor condition met")
diff --git
a/providers/influxdb/src/airflow/providers/influxdb/triggers/influxdb3.py
b/providers/influxdb/src/airflow/providers/influxdb/triggers/influxdb3.py
index f9bf999fca2..6bdeabf43e9 100644
--- a/providers/influxdb/src/airflow/providers/influxdb/triggers/influxdb3.py
+++ b/providers/influxdb/src/airflow/providers/influxdb/triggers/influxdb3.py
@@ -22,7 +22,7 @@ import asyncio
from typing import TYPE_CHECKING, Any
from airflow.providers.influxdb.hooks.influxdb3 import InfluxDB3Hook
-from airflow.providers.influxdb.utils import _convert_dataframe_to_records
+from airflow.providers.influxdb.utils import _convert_dataframe_to_records,
_first_cell_is_truthy
from airflow.triggers.base import BaseTrigger, TriggerEvent
if TYPE_CHECKING:
@@ -76,3 +76,53 @@ class InfluxDB3QueryTrigger(BaseTrigger):
return
yield TriggerEvent({"status": "success", "records": records})
+
+
+class InfluxDB3SensorTrigger(BaseTrigger):
+ """Poll an InfluxDB 3.x SQL query until its first cell meets the sensor
condition."""
+
+ def __init__(
+ self,
+ sql: str,
+ influxdb3_conn_id: str = "influxdb3_default",
+ poll_interval: float = 60,
+ fail_on_empty: bool = False,
+ ) -> None:
+ super().__init__()
+ self.sql = sql
+ self.influxdb3_conn_id = influxdb3_conn_id
+ self.poll_interval = poll_interval
+ self.fail_on_empty = fail_on_empty
+
+ def serialize(self) -> tuple[str, dict[str, Any]]:
+ return (
+
"airflow.providers.influxdb.triggers.influxdb3.InfluxDB3SensorTrigger",
+ {
+ "sql": self.sql,
+ "influxdb3_conn_id": self.influxdb3_conn_id,
+ "poll_interval": self.poll_interval,
+ "fail_on_empty": self.fail_on_empty,
+ },
+ )
+
+ async def run(self) -> AsyncIterator[TriggerEvent]:
+ hook = InfluxDB3Hook(conn_id=self.influxdb3_conn_id)
+ while True:
+ try:
+ dataframe = await hook.query_async(self.sql)
+ except Exception as error:
+ self.log.exception("InfluxDB 3 sensor query failed")
+ yield TriggerEvent({"status": "error", "message": str(error)})
+ return
+
+ if dataframe.empty and self.fail_on_empty:
+ yield TriggerEvent(
+ {"status": "fail", "message": "No rows returned, raising
as per fail_on_empty flag"}
+ )
+ return
+
+ if _first_cell_is_truthy(dataframe):
+ yield TriggerEvent({"status": "success"})
+ return
+
+ await asyncio.sleep(self.poll_interval)
diff --git a/providers/influxdb/src/airflow/providers/influxdb/utils.py
b/providers/influxdb/src/airflow/providers/influxdb/utils.py
index 79abb8d6ee2..6a64526f56d 100644
--- a/providers/influxdb/src/airflow/providers/influxdb/utils.py
+++ b/providers/influxdb/src/airflow/providers/influxdb/utils.py
@@ -28,3 +28,19 @@ if TYPE_CHECKING:
def _convert_dataframe_to_records(dataframe: pd.DataFrame) -> list[dict[str,
Any]]:
"""Convert a query result DataFrame into a JSON-serializable list of
dictionaries."""
return json.loads(dataframe.to_json(orient="records", date_format="iso"))
+
+
+def _first_cell_is_truthy(dataframe: pd.DataFrame) -> bool:
+ """Return whether the first cell meets the sensor condition."""
+ import pandas as pd
+
+ if dataframe.empty or dataframe.shape[1] == 0:
+ return False
+
+ value = dataframe.iat[0, 0]
+ if not pd.api.types.is_scalar(value):
+ raise TypeError("The first query result cell must be a scalar value")
+ if pd.isna(value):
+ return False
+
+ return value not in (0, "0", "", None)
diff --git a/providers/influxdb/tests/system/influxdb/example_influxdb3.py
b/providers/influxdb/tests/system/influxdb/example_influxdb3.py
index ef8029e58b9..c443f7df8ce 100644
--- a/providers/influxdb/tests/system/influxdb/example_influxdb3.py
+++ b/providers/influxdb/tests/system/influxdb/example_influxdb3.py
@@ -35,6 +35,7 @@ except ImportError:
from airflow.models.dag import DAG
from airflow.providers.influxdb.hooks.influxdb3 import InfluxDB3Hook
from airflow.providers.influxdb.operators.influxdb3 import InfluxDB3Operator
+from airflow.providers.influxdb.sensors.influxdb3 import InfluxDB3Sensor
@task(task_id="write_data")
@@ -66,6 +67,17 @@ deferrable_query_task = InfluxDB3Operator(
)
# [END howto_operator_influxdb3_deferrable]
+# [START howto_sensor_influxdb3]
+wait_for_data = InfluxDB3Sensor(
+ task_id="wait_for_data",
+ sql="""SELECT 1 FROM "temperature" WHERE time > now() - INTERVAL '1 hour'
LIMIT 1""",
+ influxdb3_conn_id="influxdb3_default",
+ poke_interval=60,
+ timeout=3600,
+ deferrable=True,
+)
+# [END howto_sensor_influxdb3]
+
ENV_ID = os.environ.get("SYSTEM_TESTS_ENV_ID")
DAG_ID = "influxdb3_example_dag"
@@ -77,7 +89,7 @@ with DAG(
tags=["example", "influxdb3"],
) as dag:
write_task = write_to_influxdb3()
- write_task >> [query_task, deferrable_query_task]
+ write_task >> wait_for_data >> [query_task, deferrable_query_task]
from tests_common.test_utils.watcher import watcher
diff --git a/providers/influxdb/src/airflow/providers/influxdb/utils.py
b/providers/influxdb/tests/unit/influxdb/sensors/__init__.py
similarity index 63%
copy from providers/influxdb/src/airflow/providers/influxdb/utils.py
copy to providers/influxdb/tests/unit/influxdb/sensors/__init__.py
index 79abb8d6ee2..13a83393a91 100644
--- a/providers/influxdb/src/airflow/providers/influxdb/utils.py
+++ b/providers/influxdb/tests/unit/influxdb/sensors/__init__.py
@@ -14,17 +14,3 @@
# KIND, either express or implied. See the License for the
# specific language governing permissions and limitations
# under the License.
-"""Internal helpers shared across the InfluxDB provider."""
-
-from __future__ import annotations
-
-import json
-from typing import TYPE_CHECKING, Any
-
-if TYPE_CHECKING:
- import pandas as pd
-
-
-def _convert_dataframe_to_records(dataframe: pd.DataFrame) -> list[dict[str,
Any]]:
- """Convert a query result DataFrame into a JSON-serializable list of
dictionaries."""
- return json.loads(dataframe.to_json(orient="records", date_format="iso"))
diff --git a/providers/influxdb/tests/unit/influxdb/sensors/test_influxdb3.py
b/providers/influxdb/tests/unit/influxdb/sensors/test_influxdb3.py
new file mode 100644
index 00000000000..f87d89d5fa3
--- /dev/null
+++ b/providers/influxdb/tests/unit/influxdb/sensors/test_influxdb3.py
@@ -0,0 +1,170 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+from __future__ import annotations
+
+from datetime import timedelta
+from unittest import mock
+
+import pandas as pd
+import pytest
+
+from airflow.providers.common.compat.sdk import AirflowFailException,
AirflowSensorTimeout, TaskDeferred
+from airflow.providers.influxdb.sensors.influxdb3 import InfluxDB3Sensor
+from airflow.providers.influxdb.triggers.influxdb3 import
InfluxDB3SensorTrigger
+
+SQL = """SELECT 1 FROM "events" WHERE time > now() - INTERVAL '1 hour' LIMIT
1"""
+CONN_ID = "test_influxdb3_conn"
+HOOK_PATH = "airflow.providers.influxdb.sensors.influxdb3.InfluxDB3Hook"
+
+
+class TestInfluxDB3Sensor:
+ def test_init(self):
+ sensor = InfluxDB3Sensor(task_id="wait", sql=SQL)
+
+ assert sensor.sql == SQL
+ assert sensor.influxdb3_conn_id == "influxdb3_default"
+ assert sensor.fail_on_empty is False
+ assert sensor.deferrable is False
+ assert sensor.template_fields == ("sql", "influxdb3_conn_id")
+ assert sensor.template_ext == (".sql",)
+
+ @pytest.mark.parametrize(
+ ("dataframe", "expected"),
+ [
+ pytest.param(pd.DataFrame({"literal": [1]}), True,
id="numeric-one"),
+ pytest.param(pd.DataFrame({"literal": ["ready"]}), True,
id="non-empty-string"),
+ pytest.param(pd.DataFrame({"literal": []}), False, id="no-rows"),
+ pytest.param(pd.DataFrame({"count": [0]}), False,
id="numeric-zero"),
+ pytest.param(pd.DataFrame({"count": ["0"]}), False,
id="string-zero"),
+ pytest.param(pd.DataFrame({"value": [None]}), False, id="none"),
+ ],
+ )
+ @mock.patch(HOOK_PATH, autospec=True)
+ def test_poke(self, mock_hook_class, dataframe, expected):
+ mock_hook_class.return_value.query.return_value = dataframe
+ sensor = InfluxDB3Sensor(task_id="wait", sql=SQL,
influxdb3_conn_id=CONN_ID)
+
+ assert sensor.poke(context={}) is expected
+ mock_hook_class.assert_called_once_with(conn_id=CONN_ID)
+ mock_hook_class.return_value.query.assert_called_once_with(SQL)
+
+ @mock.patch(HOOK_PATH, autospec=True)
+ def test_poke_fail_on_empty(self, mock_hook_class):
+ mock_hook_class.return_value.query.return_value =
pd.DataFrame({"literal": []})
+ sensor = InfluxDB3Sensor(task_id="wait", sql=SQL, fail_on_empty=True)
+
+ with pytest.raises(AirflowFailException, match="fail_on_empty"):
+ sensor.poke(context={})
+
+ @mock.patch.object(InfluxDB3Sensor, "defer", autospec=True)
+ @mock.patch(HOOK_PATH, autospec=True)
+ def test_execute_deferrable_fail_on_empty_before_deferral(self,
mock_hook_class, mock_defer):
+ mock_hook_class.return_value.query.return_value =
pd.DataFrame({"literal": []})
+ sensor = InfluxDB3Sensor(
+ task_id="wait",
+ sql=SQL,
+ fail_on_empty=True,
+ deferrable=True,
+ )
+
+ with pytest.raises(AirflowFailException, match="fail_on_empty"):
+ sensor.execute(context={})
+
+ mock_hook_class.return_value.query.assert_called_once_with(SQL)
+ mock_defer.assert_not_called()
+
+ @mock.patch(HOOK_PATH, autospec=True)
+ def test_execute_times_out_when_condition_is_never_met(self,
mock_hook_class):
+ mock_hook_class.return_value.query.return_value =
pd.DataFrame({"literal": []})
+ sensor = InfluxDB3Sensor(task_id="wait", sql=SQL, poke_interval=0,
timeout=0)
+
+ with pytest.raises(AirflowSensorTimeout):
+ sensor.execute(context={})
+
+ @mock.patch(HOOK_PATH, autospec=True)
+ def test_execute_deferrable_completes_after_initial_match(self,
mock_hook_class):
+ mock_hook_class.return_value.query.return_value =
pd.DataFrame({"literal": [1]})
+ sensor = InfluxDB3Sensor(task_id="wait", sql=SQL, deferrable=True)
+
+ assert sensor.execute(context={}) is None
+ mock_hook_class.return_value.query.assert_called_once_with(SQL)
+
+ @mock.patch(HOOK_PATH, autospec=True)
+ def test_execute_deferrable_defers_with_sensor_settings(self,
mock_hook_class):
+ mock_hook_class.return_value.query.return_value =
pd.DataFrame({"literal": []})
+ sensor = InfluxDB3Sensor(
+ task_id="wait",
+ sql=SQL,
+ influxdb3_conn_id=CONN_ID,
+ fail_on_empty=False,
+ deferrable=True,
+ poke_interval=30,
+ timeout=600,
+ )
+
+ with pytest.raises(TaskDeferred) as exc:
+ sensor.execute(context={})
+
+ trigger = exc.value.trigger
+ assert isinstance(trigger, InfluxDB3SensorTrigger)
+ assert trigger.sql == SQL
+ assert trigger.influxdb3_conn_id == CONN_ID
+ assert trigger.poll_interval == 30
+ assert trigger.fail_on_empty is False
+ assert exc.value.method_name == "execute_complete"
+ assert exc.value.timeout == timedelta(minutes=10)
+
+ def test_execute_complete_success(self):
+ sensor = InfluxDB3Sensor(task_id="wait", sql=SQL, deferrable=True)
+
+ assert sensor.execute_complete(context={}, event={"status":
"success"}) is None
+
+ @pytest.mark.parametrize(
+ ("event", "match"),
+ [
+ pytest.param(
+ {"status": "fail", "message": "No rows returned, raising as
per fail_on_empty flag"},
+ "fail_on_empty",
+ id="fail-with-message",
+ ),
+ pytest.param(
+ {"status": "fail"},
+ "InfluxDB 3 sensor failed",
+ id="fail-without-message",
+ ),
+ ],
+ )
+ def test_execute_complete_fail_on_empty(self, event, match):
+ sensor = InfluxDB3Sensor(task_id="wait", sql=SQL, deferrable=True)
+
+ with pytest.raises(AirflowFailException, match=match):
+ sensor.execute_complete(context={}, event=event)
+
+ @pytest.mark.parametrize(
+ ("event", "match"),
+ [
+ pytest.param(None, "did not return an event", id="missing-event"),
+ pytest.param({"status": "error", "message": "boom"}, "boom",
id="error-with-message"),
+ pytest.param({"status": "error"}, "InfluxDB 3 sensor failed",
id="error-without-message"),
+ pytest.param({"status": "cancelled"}, "unexpected status",
id="unexpected-status"),
+ ],
+ )
+ def test_execute_complete_failures(self, event, match):
+ sensor = InfluxDB3Sensor(task_id="wait", sql=SQL, deferrable=True)
+
+ with pytest.raises(RuntimeError, match=match):
+ sensor.execute_complete(context={}, event=event)
diff --git a/providers/influxdb/tests/unit/influxdb/test_utils.py
b/providers/influxdb/tests/unit/influxdb/test_utils.py
index 270474453b1..1e59ea7b639 100644
--- a/providers/influxdb/tests/unit/influxdb/test_utils.py
+++ b/providers/influxdb/tests/unit/influxdb/test_utils.py
@@ -17,8 +17,9 @@
from __future__ import annotations
import pandas as pd
+import pytest
-from airflow.providers.influxdb.utils import _convert_dataframe_to_records
+from airflow.providers.influxdb.utils import _convert_dataframe_to_records,
_first_cell_is_truthy
def test_convert_dataframe_to_records_serializes_rows_and_timestamps():
@@ -33,3 +34,29 @@ def
test_convert_dataframe_to_records_serializes_rows_and_timestamps():
{"col1": 1, "timestamp": "2024-01-01T00:00:00.000Z"},
{"col1": 2, "timestamp": "2024-01-02T03:04:05.000Z"},
]
+
+
[email protected](
+ ("dataframe", "expected"),
+ [
+ pytest.param(pd.DataFrame({"literal": [1]}), True, id="numeric-one"),
+ pytest.param(pd.DataFrame({"count": [42]}), True, id="positive-count"),
+ pytest.param(pd.DataFrame({"flag": ["ready"]}), True,
id="non-empty-string"),
+ pytest.param(pd.DataFrame({"literal": []}), False, id="no-rows"),
+ pytest.param(pd.DataFrame(), False, id="no-columns"),
+ pytest.param(pd.DataFrame({"count": [0]}), False, id="numeric-zero"),
+ pytest.param(pd.DataFrame({"count": ["0"]}), False, id="string-zero"),
+ pytest.param(pd.DataFrame({"value": [float("nan")]}), False, id="nan"),
+ pytest.param(pd.DataFrame({"value": [None]}), False, id="none"),
+ pytest.param(pd.DataFrame({"first": [0, 1], "second": [1, 1]}), False,
id="first-cell-only"),
+ ],
+)
+def test_first_cell_is_truthy(dataframe, expected):
+ assert _first_cell_is_truthy(dataframe) is expected
+
+
+def test_first_cell_is_truthy_rejects_non_scalar_value():
+ dataframe = pd.DataFrame({"value": [[1, 2]]})
+
+ with pytest.raises(TypeError, match="must be a scalar"):
+ _first_cell_is_truthy(dataframe)
diff --git a/providers/influxdb/tests/unit/influxdb/triggers/test_influxdb3.py
b/providers/influxdb/tests/unit/influxdb/triggers/test_influxdb3.py
index 18669c504cf..89de9cf467c 100644
--- a/providers/influxdb/tests/unit/influxdb/triggers/test_influxdb3.py
+++ b/providers/influxdb/tests/unit/influxdb/triggers/test_influxdb3.py
@@ -22,7 +22,7 @@ from unittest import mock
import pandas as pd
import pytest
-from airflow.providers.influxdb.triggers.influxdb3 import InfluxDB3QueryTrigger
+from airflow.providers.influxdb.triggers.influxdb3 import
InfluxDB3QueryTrigger, InfluxDB3SensorTrigger
from airflow.triggers.base import TriggerEvent
SQL = 'SELECT "duration" FROM "pyexample"'
@@ -79,3 +79,124 @@ class TestInfluxDB3QueryTrigger:
with pytest.raises(asyncio.CancelledError):
await anext(trigger.run())
+
+
+class TestInfluxDB3SensorTrigger:
+ def test_serialization(self):
+ """Trigger serializes its constructor arguments."""
+ trigger = InfluxDB3SensorTrigger(
+ sql=SQL,
+ influxdb3_conn_id=CONN_ID,
+ poll_interval=30,
+ fail_on_empty=True,
+ )
+ classpath, kwargs = trigger.serialize()
+
+ assert classpath ==
"airflow.providers.influxdb.triggers.influxdb3.InfluxDB3SensorTrigger"
+ assert kwargs == {
+ "sql": SQL,
+ "influxdb3_conn_id": CONN_ID,
+ "poll_interval": 30,
+ "fail_on_empty": True,
+ }
+
+ @pytest.mark.asyncio
+ @mock.patch("airflow.providers.influxdb.triggers.influxdb3.asyncio.sleep",
new_callable=mock.AsyncMock)
+ @mock.patch("airflow.providers.influxdb.triggers.influxdb3.InfluxDB3Hook",
autospec=True)
+ async def test_run_succeeds_without_sleep_when_condition_is_met(self,
mock_hook_class, mock_sleep):
+ mock_hook = mock_hook_class.return_value
+ mock_hook.query_async =
mock.AsyncMock(return_value=pd.DataFrame({"literal": [1]}))
+
+ events = [event async for event in
InfluxDB3SensorTrigger(sql=SQL).run()]
+
+ mock_hook_class.assert_called_once_with(conn_id="influxdb3_default")
+ mock_hook.query_async.assert_awaited_once_with(SQL)
+ mock_sleep.assert_not_awaited()
+ assert events == [TriggerEvent({"status": "success"})]
+
+ @pytest.mark.asyncio
+ @mock.patch("airflow.providers.influxdb.triggers.influxdb3.asyncio.sleep",
new_callable=mock.AsyncMock)
+ @mock.patch("airflow.providers.influxdb.triggers.influxdb3.InfluxDB3Hook",
autospec=True)
+ async def test_run_polls_until_condition_is_met(self, mock_hook_class,
mock_sleep):
+ mock_hook = mock_hook_class.return_value
+ mock_hook.query_async = mock.AsyncMock(
+ side_effect=[
+ pd.DataFrame({"literal": []}),
+ pd.DataFrame({"count": [0]}),
+ pd.DataFrame({"literal": [1]}),
+ ]
+ )
+
+ events = [event async for event in InfluxDB3SensorTrigger(sql=SQL,
poll_interval=30).run()]
+
+ assert mock_hook.query_async.await_count == 3
+ assert mock_sleep.await_args_list == [mock.call(30), mock.call(30)]
+ assert events == [TriggerEvent({"status": "success"})]
+
+ @pytest.mark.asyncio
+ @mock.patch("airflow.providers.influxdb.triggers.influxdb3.asyncio.sleep",
new_callable=mock.AsyncMock)
+ @mock.patch("airflow.providers.influxdb.triggers.influxdb3.InfluxDB3Hook",
autospec=True)
+ async def test_run_fail_on_empty(self, mock_hook_class, mock_sleep):
+ mock_hook = mock_hook_class.return_value
+ mock_hook.query_async =
mock.AsyncMock(return_value=pd.DataFrame({"literal": []}))
+
+ events = [event async for event in InfluxDB3SensorTrigger(sql=SQL,
fail_on_empty=True).run()]
+
+ mock_hook.query_async.assert_awaited_once_with(SQL)
+ mock_sleep.assert_not_awaited()
+ assert events == [
+ TriggerEvent({"status": "fail", "message": "No rows returned,
raising as per fail_on_empty flag"})
+ ]
+
+ @pytest.mark.asyncio
+ @mock.patch("airflow.providers.influxdb.triggers.influxdb3.InfluxDB3Hook",
autospec=True)
+ async def test_run_failure(self, mock_hook_class):
+ mock_hook = mock_hook_class.return_value
+ mock_hook.query_async = mock.AsyncMock(side_effect=ValueError("boom"))
+
+ events = [event async for event in
InfluxDB3SensorTrigger(sql=SQL).run()]
+
+ mock_hook_class.assert_called_once_with(conn_id="influxdb3_default")
+ mock_hook.query_async.assert_awaited_once_with(SQL)
+ assert events == [TriggerEvent({"status": "error", "message": "boom"})]
+
+ @pytest.mark.asyncio
+ @mock.patch("airflow.providers.influxdb.triggers.influxdb3.asyncio.sleep",
new_callable=mock.AsyncMock)
+ @mock.patch("airflow.providers.influxdb.triggers.influxdb3.InfluxDB3Hook",
autospec=True)
+ async def test_run_failure_after_unsuccessful_poll(self, mock_hook_class,
mock_sleep):
+ mock_hook = mock_hook_class.return_value
+ mock_hook.query_async = mock.AsyncMock(
+ side_effect=[pd.DataFrame({"literal": []}), ValueError("boom")]
+ )
+
+ events = [event async for event in InfluxDB3SensorTrigger(sql=SQL,
poll_interval=30).run()]
+
+ assert mock_hook.query_async.await_count == 2
+ mock_sleep.assert_awaited_once_with(30)
+ assert events == [TriggerEvent({"status": "error", "message": "boom"})]
+
+ @pytest.mark.asyncio
+ @mock.patch("airflow.providers.influxdb.triggers.influxdb3.InfluxDB3Hook",
autospec=True)
+ async def test_run_propagates_cancellation_during_query(self,
mock_hook_class):
+ mock_hook = mock_hook_class.return_value
+ mock_hook.query_async =
mock.AsyncMock(side_effect=asyncio.CancelledError())
+
+ with pytest.raises(asyncio.CancelledError):
+ await anext(InfluxDB3SensorTrigger(sql=SQL).run())
+
+ @pytest.mark.asyncio
+ @mock.patch(
+ "airflow.providers.influxdb.triggers.influxdb3.asyncio.sleep",
+ new_callable=mock.AsyncMock,
+ side_effect=asyncio.CancelledError(),
+ )
+ @mock.patch("airflow.providers.influxdb.triggers.influxdb3.InfluxDB3Hook",
autospec=True)
+ async def test_run_propagates_cancellation_during_sleep(self,
mock_hook_class, mock_sleep):
+ mock_hook = mock_hook_class.return_value
+ mock_hook.query_async =
mock.AsyncMock(return_value=pd.DataFrame({"literal": []}))
+
+ with pytest.raises(asyncio.CancelledError):
+ await anext(InfluxDB3SensorTrigger(sql=SQL,
poll_interval=30).run())
+
+ mock_hook.query_async.assert_awaited_once_with(SQL)
+ mock_sleep.assert_awaited_once_with(30)