Skip to content

Add dedicated InfluxDB 3 sensor - #73522

Merged
shahar1 merged 12 commits into
apache:mainfrom
subhramit:influxdb3-sensor
Sep 25, 2026
Merged

shahar1 merged 12 commits into
apache:mainfrom
subhramit:influxdb3-sensor

Conversation

@subhramit

@subhramit subhramit commented Sep 22, 2026 •

Copy link
Copy Markdown
Contributor

Implement Design A from #67109 by adding a generic SQL-truthy InfluxDB3Sensor.

Now that deferrable mode is available (#71976), implementing the sensor felt like the obvious next step.
The sensor evaluates the first cell returned by an InfluxDB 3 SQL query, following SqlSensor-style semantics. It supports synchronous polling and deferrable polling through a new InfluxDB3SensorTrigger, which reruns the query at poke_interval without occupying a worker slot between checks. Airflow owns the sensor timeout, while query errors and trigger cancellation retain their existing failure and cancellation behavior.

Includes a deferrable system-test example. Tests cover timeout propagation and cancellation during both query execution and sleep.

Have locally run all prek hooks and tests and all of them pass. Following is the manual test.

MWE

Tested against influxdb:3-core with a file object store. Two tables in airflow_demo: home for the client check, and arrivals seeded with a single row that fails the sensor's predicate - so the table exists but the query comes back empty and the sensor has something real to poll.

image

The async client on its own against the same server, to isolate it from Airflow:

c = InfluxDBClient3(host="http://localhost:8181", token=TOKEN, database="airflow_demo")
print(asyncio.run(c.query_async("SELECT * FROM home")).to_pydict())
image

Then through the sensor - deferrable, gating a downstream operator so the blocking behaviour is visible:

wait_for_data = InfluxDB3Sensor(
    task_id="wait_for_data",
    sql='SELECT 1 FROM "arrivals" WHERE ready = 1 LIMIT 1',
    influxdb3_conn_id=CONN,
    deferrable=True,
    poke_interval=15,
    timeout=300,
)
read_data = InfluxDB3Operator(task_id="read_data", sql='SELECT * FROM "arrivals" WHERE ready = 1', influxdb3_conn_id=CONN)
wait_for_data >> read_data

Triggered and left to poll - wait_for_data sits deferred while read_data has no state at all:

image

Then the matching row, written by hand:

image

The initial poke came back empty, so the task deferred and the trigger polled every 15s. The poke after the write caught the new row, fired TriggerEvent<{'status': 'success'}>, and the task resumed - total duration 00:02:52, waiting on the data:

image

Gantt view of the same run (showing the sensor's wait dominating the run and read_data running only after it cleared):

image

Both tasks finally from the CLI:
image

read_data's XCom shows it read the row the sensor was waiting for. Note the seeded ready=0 row is correctly excluded:

image

Also note: the URI: line repeats per poke because the trigger calls get_conn() on each loop iteration (which is expected).


Was generative AI tooling used to co-author this PR?
  • Yes

Assisted-by: Zed GPT-5.6 Sol following the guidelines

Note: All code changes done as a result (and also this description) were manually driven, edited & reviewed by me.

@subhramit subhramit changed the title Add deferrable InfluxDB 3 sensor Add dedicated InfluxDB 3 sensor Sep 22, 2026
@subhramit

subhramit commented Sep 22, 2026 •

Copy link
Copy Markdown
Contributor Author

Adding reviewers of the previous PR
cc @eladkal, @SameerMesiah97
@ashb @potiuk

Signed-off-by: subhramit <subhramit.bb@live.in>
Signed-off-by: Subhramit Basu <subhramit.bb@live.in>
Signed-off-by: subhramit <subhramit.bb@live.in>

@SameerMesiah97 SameerMesiah97 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Looks good. There are 2 additional optional scenarios you can cover in the tests. I have left some comments.

Comment thread providers/influxdb/src/airflow/providers/influxdb/sensors/influxdb3.py Outdated
Comment thread providers/influxdb/src/airflow/providers/influxdb/sensors/influxdb3.py Outdated
Comment thread providers/influxdb/src/airflow/providers/influxdb/utils.py
Comment thread providers/influxdb/tests/unit/influxdb/triggers/test_influxdb3.py Outdated
Comment thread providers/influxdb/tests/unit/influxdb/sensors/test_influxdb3.py
Comment thread providers/influxdb/tests/unit/influxdb/triggers/test_influxdb3.py
Signed-off-by: subhramit <subhramit.bb@live.in>
Signed-off-by: subhramit <subhramit.bb@live.in>
Signed-off-by: subhramit <subhramit.bb@live.in>
Signed-off-by: subhramit <subhramit.bb@live.in>
Signed-off-by: subhramit <subhramit.bb@live.in>

@justinpakzad justinpakzad left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the PR. I left a couple of comments.

Comment thread providers/influxdb/src/airflow/providers/influxdb/utils.py Outdated
Signed-off-by: subhramit <subhramit.bb@live.in>
Signed-off-by: subhramit <subhramit.bb@live.in>

@SameerMesiah97 SameerMesiah97 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I have left 2 more comments for you to address before merge.

Signed-off-by: subhramit <subhramit.bb@live.in>

@aaron-y-chen aaron-y-chen left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM! a nit, not a blocker :)

Signed-off-by: subhramit <subhramit.bb@live.in>
@shahar1
shahar1 merged commit c1e98dc into apache:main Sep 25, 2026
80 checks passed
@subhramit
subhramit deleted the influxdb3-sensor branch September 25, 2026 07:20
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

InfluxDB 3: add a dedicated sensor (data freshness / windowed existence)

5 participants