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

list_tables

Lists registered table names

get_schema

Returns a table’s columns (Arrow schema) as JSON, with a name_contains filter and a bounded summary on very wide tables (see Bounded schema results)

query

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:

airflow/providers/common/ai/example_dags/example_datafusion_toolset.py[source]

    @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 payroll

error: 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 sales

error: 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 more DataSourceConfig entries. Requires apache-airflow-providers-common-sql[datafusion].

  • allow_writes: Allow data-modifying SQL (CREATE TABLE, CREATE VIEW, INSERT INTO, etc.). Default False: 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 the query tool. Default 50.

  • max_result_bytes: Budget for the serialized query result, and the byte backstop that also triggers the get_schema summary. Default 64 KiB. See Bounded query results and Bounded schema results.

  • max_columns: Maximum columns get_schema returns in full. Default 100. 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. Default None, the agent’s retries. 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=False is 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 query tool runs the statement with a LIMIT of max_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.

Was this entry helpful?