Skip to content
Merged
1 change: 1 addition & 0 deletions providers/influxdb/docs/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
40 changes: 40 additions & 0 deletions providers/influxdb/docs/sensors/index.rst
Original file line number Diff line number Diff line change
@@ -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.
5 changes: 5 additions & 0 deletions providers/influxdb/provider.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,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:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
# 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.
108 changes: 108 additions & 0 deletions providers/influxdb/src/airflow/providers/influxdb/sensors/influxdb3.py
Original file line number Diff line number Diff line change
@@ -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.

Comment thread
shahar1 marked this conversation as resolved.
.. 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")
Comment thread
subhramit marked this conversation as resolved.
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}")
Comment thread
subhramit marked this conversation as resolved.

self.log.info("InfluxDB 3 sensor condition met")
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@
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:
Expand Down Expand Up @@ -76,3 +76,53 @@ async def run(self) -> AsyncIterator[TriggerEvent]:
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)
16 changes: 16 additions & 0 deletions providers/influxdb/src/airflow/providers/influxdb/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,3 +28,19 @@
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
Comment thread
subhramit marked this conversation as resolved.

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)
14 changes: 13 additions & 1 deletion providers/influxdb/tests/system/influxdb/example_influxdb3.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@
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")
Expand Down Expand Up @@ -66,6 +67,17 @@ def write_to_influxdb3():
)
# [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"

Expand All @@ -77,7 +89,7 @@ def write_to_influxdb3():
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

Expand Down
16 changes: 16 additions & 0 deletions providers/influxdb/tests/unit/influxdb/sensors/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
# 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.
Loading