justinpakzad commented on code in PR #73522: URL: https://github.com/apache/airflow/pull/73522#discussion_r4088140844
########## providers/influxdb/src/airflow/providers/influxdb/sensors/influxdb3.py: ########## @@ -0,0 +1,102 @@ +# 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. + + :param sql: The SQL query to poll. + :param influxdb3_conn_id: Reference to :ref:`InfluxDB 3 connection id <howto/connection:influxdb3>`. + :param fail_on_empty: Fail instead of waiting when the query returns no rows. + :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") Review Comment: There are some differences in behavior with fail on empty between the sync and deferrable path. The sync path here raises `AirflowFailException`, which will ignore any remaining retry attempts (as per the docs), but the deferrable path would raise a `RuntimeError` in `execute_complete`, which honors the retry policy. I think the two should be aligned? ########## providers/influxdb/src/airflow/providers/influxdb/utils.py: ########## @@ -28,3 +28,20 @@ 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] + try: + if bool(pd.isna(value)): + return False + except ValueError as error: + raise TypeError("The first query result cell must be a scalar value") from error Review Comment: This is optional but instead of catching a `ValueError` and re-raising a `TypeError`, this could be simplified to: ``` 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 -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
