InfluxDBOperator

Use the InfluxDBOperator to execute SQL commands in a InfluxDB database.

An example of running the query using the operator:

tests/system/influxdb/example_influxdb_query.py[source]


    query_influxdb_task = InfluxDBOperator(
        influxdb_conn_id="influxdb_conn_id",
        task_id="query_influxdb",
        sql='from(bucket:"test-influx") |> range(start: -10m, stop: {{ds}})',
        dag=dag,
    )

InfluxDB3Operator

The InfluxDB3Operator executes SQL queries in an InfluxDB 3.x database.

Example usage:

tests/system/influxdb/example_influxdb3.py[source]

query_task = InfluxDB3Operator(
    task_id="query_data",
    sql="SELECT * FROM \"temperature\" WHERE time > now() - INTERVAL '1 hour'",
    influxdb3_conn_id="influxdb3_default",
)

Deferrable mode

Set deferrable=True to release the worker slot while the query runs. The task is resumed by the InfluxDB3QueryTrigger once results are ready.

tests/system/influxdb/example_influxdb3.py[source]

deferrable_query_task = InfluxDB3Operator(
    task_id="query_data_deferrable",
    sql="""SELECT * FROM "temperature" WHERE time > now() - INTERVAL '1 hour'""",
    influxdb3_conn_id="influxdb3_default",
    deferrable=True,
)

Note

This implementation follows the upstream influxdb3-python client, which documents querying through the Flight client and exposes query_async() for asynchronous query execution.

Results still travel back through XCom, so deferring is most useful for long-running queries with small-to-moderate result sets. For very large extracts, keep using InfluxDB3Hook from a Python task.

Was this entry helpful?