InfluxDBOperator¶
Use the InfluxDBOperator to execute
SQL commands in a InfluxDB database.
An example of running the query using the operator:
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:
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.
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.