Files with DataFusion: DataFusionToolset¶
Curated toolset wrapping
DataFusionEngine
with three tools (list_tables, get_schema, and query) for
querying files on object stores (S3, GCS, local filesystem, Iceberg) via Apache DataFusion.
Tool |
Description |
|---|---|
|
Lists registered table names |
|
Returns column names and types for a table (Arrow schema) |
|
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.
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. Default 64 KiB. See Bounded query results.
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.
A real example. The same bucket as the HookToolset example on Airflow hooks as tools: HookToolset, reached
the other way. Rather than exposing list_keys and read_key and leaving
the agent to reassemble files, this registers the prefix as a table and lets it
write SQL:
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",
),
],
max_rows=100,
)
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.