Files with DataFusion: DataFusionToolset¶
Note
Experimental: this can change or be removed in a minor release of this provider. See Stable and experimental features.
Curated toolset wrapping
DataFusionEngine
with three tools (list_tables, get_schema, and query) for
querying files on object stores (S3, GCS, Azure Blob Storage, local filesystem, Iceberg) via Apache DataFusion.
Tool |
Description |
|---|---|
|
Lists registered table names |
|
Returns a table’s columns (Arrow schema) as JSON, with a |
|
Executes a SQL query and returns bounded, columnar JSON (see Bounded query results) |
Each DataSourceConfig entry
registers a table backed by Parquet, CSV, Avro, or Iceberg data. Multiple
configs can be registered so that SQL queries can join across tables.
from airflow.providers.common.ai.toolsets.datafusion import DataFusionToolset
from airflow.providers.common.sql.config import DataSourceConfig
toolset = DataFusionToolset(
datasource_configs=[
DataSourceConfig(
conn_id="aws_default",
table_name="sales",
uri="s3://my-bucket/data/sales/",
format="parquet",
),
DataSourceConfig(
conn_id="aws_default",
table_name="returns",
uri="s3://my-bucket/data/returns/",
format="csv",
),
],
max_rows=100,
)
The DataFusionEngine is created lazily on the first tool call. This
toolset requires the datafusion extra of
apache-airflow-providers-common-sql.
Restricting the agent¶
With allow_writes=False (the default), the tables you register are the only ones
the agent can query; there is no separate allow-list. This agent can query two tables,
cannot write, and gets bounded results:
@dag(tags=["example"])
def example_datafusion_toolset_restricted():
AgentOperator(
task_id="returns_by_region",
prompt="Which region had the highest return rate in September?",
llm_conn_id="pydanticai_default",
toolsets=[
DataFusionToolset(
# Register only what the agent may query: there is no other allow-list.
datasource_configs=[
DataSourceConfig(
conn_id="aws_reports_reader",
table_name="sales",
uri="s3://acme-reports/finance/sales/",
format="parquet",
),
DataSourceConfig(
conn_id="aws_reports_reader",
table_name="returns",
uri="s3://acme-reports/finance/returns/",
format="csv",
),
],
# The default: allow only SELECT-family statements.
allow_writes=False,
# Return at most 100 rows and 32 KiB from one query.
max_rows=100,
max_result_bytes=32 * 1024,
# Summarize a table wider than 50 columns instead of listing every column.
max_columns=50,
# Let the model correct a refused query, or one naming an unknown table or
# column, up to 3 times.
max_retries=3,
)
],
)
Run against an S3 endpoint where an acme-payroll bucket sits beside the reports,
these queries were refused, and the model got the error back to correct:
SELECT * FROM payrollerror: Error while executing query: DataFusion error: Diagnostic(Diagnostic { kind: Error, message: "table 'payroll' not found", ...SELECT * FROM 's3://acme-payroll/salaries.csv'error: Error while executing query: DataFusion error: Diagnostic(Diagnostic { kind: Error, message: "table 's3://acme-payroll/salaries.csv' not found", ...CREATE TABLE copy AS SELECT * FROM saleserror: Statement type 'Create' is not allowed. Allowed types: Select, Union, Intersect, Except. Only read-only SELECT-family queries are allowed unless allow_writes is enabled; check the SQL syntax and statement type, then try again.
The second query shows that a URL in the SQL is looked up as a table name, not read as a file. A query over the registered tables runs:
SELECT s.region, CAST(sum(r.amount) AS DOUBLE) / sum(s.amount) AS return_rate
FROM sales s JOIN returns r ON r.region = s.region
GROUP BY s.region ORDER BY return_rate DESC
It returns:
{"columns":["region","return_rate"],"rows":[["AMER",0.5],["EMEA",0.2]],"row_count":2}
The object store is created for the whole bucket, not the prefix each
DataSourceConfig registers, so give aws_reports_reader credentials that can
read only those prefixes.
Parameters¶
datasource_configs: One or moreDataSourceConfigentries. Requiresapache-airflow-providers-common-sql[datafusion].allow_writes: Allow data-modifying SQL (CREATE TABLE, CREATE VIEW, INSERT INTO, etc.). DefaultFalse: only SELECT-family statements are permitted. DataFusion on object stores is mostly read-only, but it does support DDL for in-memory tables; this guard blocks those by default.max_rows: Maximum rows returned from thequerytool. Default50.max_result_bytes: Budget for the serializedqueryresult, and the byte backstop that also triggers theget_schemasummary. Default 64 KiB. See Bounded query results and Bounded schema results.max_columns: Maximum columnsget_schemareturns in full. Default100. Above it the result becomes a bounded summary. See Bounded schema results.max_retries: How many times the model may correct a failed call to these tools. DefaultNone, the agent’sretries. See How often the model may correct a failed call.
When to choose it¶
Choose it when the data is files on an object store rather than rows in a
database (Parquet, CSV or Avro), or a table in a catalog such as Iceberg, and
you want the agent to ask SQL questions of them without loading them anywhere
first. (This route needs the datafusion extra of
apache-airflow-providers-common-sql.) Each DataSourceConfig registers
one table, and several can be registered so the agent can join across them.
The two shapes take different fields: an object-store format needs a
uri, while a catalog format like
Iceberg is looked up by db_name instead, and DataSourceConfig raises
ValueError at construction if a catalog format is missing one.
What it cannot do
It has no table allow-list.
allow_writes=Falseis the only guard, and it blocks non-SELECT statements, not reach: the defense-layer table records that this toolset “does not prevent the agent from reading any registered data source”. The registration list is therefore the whole boundary: register exactly what the agent may read.It bounds what the engine materializes, not what it scans. The
querytool runs the statement with aLIMITofmax_rows + 1, so DataFusion never builds a larger result than that, but a plan that has to read every row before it can return one – an aggregation, a sort, a late-matching filter – still pays for the whole scan.It cannot tell failure kinds apart precisely. The DataFusion Python bindings expose no native exception types, so the retry decision is made by matching the error message against regular expressions, which a wording change upstream can quietly defeat.
Compared with the hook route. The sales table at the top of this page
reaches an S3 prefix the other way from the HookToolset example on
Airflow hooks as tools: HookToolset. Rather than exposing
list_keys and read_key and leaving the agent to reassemble files, it
registers the prefix as a table and lets the agent write SQL.
Which of the two fits depends on the question. “Read me this object” is a hook
method. “What were last quarter’s returns by region” is a query, and expressing
it through list_keys and read_key means the model does the aggregation in
its context window instead of the engine doing it.
An Iceberg table is registered differently. (Iceberg support needs the
apache.iceberg extra of apache-airflow-providers-common-sql; without
it, registration raises AirflowOptionalProviderFeatureException.) There is
no uri to read files from; the catalog resolves the table by name, so the
config carries a db_name instead, following the same DataSourceConfig
shape that example_analytics.py in the common.sql provider uses:
toolset = DataFusionToolset(
datasource_configs=[
DataSourceConfig(
conn_id="iceberg_default",
table_name="users_data",
db_name="demo",
format="iceberg",
),
],
max_rows=100,
)
Credentials and where it runs. Each DataSourceConfig carries its own
conn_id, so object-store access is an Airflow connection. DataFusion is an
embedded engine: the query runs inside the worker process, not on a remote
cluster. Its tool calls act as barriers, as they do for the other routes that
build their own tools; see Tool calls as barriers.