airflow.providers.common.ai.toolsets.datafusion¶
Curated SQL toolset wrapping DataFusionEngine for agentic object-store workflows.
Attributes¶
Classes¶
Curated toolset that gives an LLM agent SQL access to object-storage data via Apache DataFusion. |
Module Contents¶
- class airflow.providers.common.ai.toolsets.datafusion.DataFusionToolset(datasource_configs, *, allow_writes=False, max_rows=50, max_result_bytes=DEFAULT_MAX_RESULT_BYTES, max_columns=DEFAULT_MAX_COLUMNS, max_retries=None)[source]¶
Bases:
airflow.providers.common.ai.utils.toolset_base.AirflowToolsetCurated toolset that gives an LLM agent SQL access to object-storage data via Apache DataFusion.
Note
Experimental: this can change or be removed in a minor release of this provider. See Stable and experimental features.
Provides three tools —
list_tables,get_schema, andquery— backed byDataFusionEngine.Each
DataSourceConfigentry registers a table backed by Parquet, CSV, Avro, or Iceberg data on S3 or local storage. Multiple configs can be registered so that SQL queries can join across tables.Requires the
datafusionextra ofapache-airflow-providers-common-sql.- Parameters:
datasource_configs (list[airflow.providers.common.sql.config.DataSourceConfig]) – One or more DataFusion data-source configurations.
allow_writes (bool) – Allow data-modifying SQL (CREATE TABLE, CREATE VIEW, INSERT INTO, etc.). Default
False— only SELECT-family statements are permitted.EXPLAINreaches the engine only withallow_writes=True, and fails there: themax_rowslimit wraps the plan, and DataFusion requiresEXPLAINto be the root of the plan. The agent gets an error result, not the plan.max_rows (int) – Maximum number of rows returned from the
querytool. Default50. The query is limited tomax_rows + 1rows, so a large result is never fully materialized; the extra row only signals truncation.max_result_bytes (int) – Budget for the serialized
queryresult, in bytes, and the byte backstop that also triggers theget_schemasummary (seemax_columns). Default 64 KiB.max_rowsbounds rows, which says nothing about size: one row of a 3000-column table is larger than a thousand rows of a narrow one, and a tool result stays in the model’s message history for the rest of the run, so its cost is re-paid on every subsequent request. Rows are returned as a contiguous prefix, stopping at the first that does not fit the remaining budget rather than skipping it and packing later ones, so one wide row early in the result ends it. The result reports which limit it hit so the agent can narrow its projection rather than page through the table.max_columns (int) – Maximum number of columns
get_schemareturns in full. Default100. Above it – or when the serialized columns exceedmax_result_bytes– the full list is replaced by a bounded summary (column count, a type histogram, a sample of columns) that points the agent at thename_containsfilter, so a several-thousand-column table cannot exhaust the context before a query is written.max_retries (int | None) – How many times the model may correct a failed call to one of these tools before the run fails.
None(the default) uses the agent’s tool retry budget, itsretries, as pydantic-ai’s own toolsets do.
- property id: str[source]¶
An ID for the toolset that is unique among all toolsets registered with the same agent.
If you’re implementing a concrete implementation that users can instantiate more than once, you should let them optionally pass a custom ID to the constructor and return that here.
A toolset needs to have an ID in order to be used in a durable execution environment like Temporal, in which case the ID will be used to identify the toolset’s activities within the workflow.
IDs wrapped in angle brackets (‘<agent>’ for an agent’s own function toolset, ‘<output>’ for its output tools) name a role the framework fills on the user’s behalf rather than a registered toolset. Don’t return one from your own toolset.
- async execute_tool(name, tool_args, *, ctx, tool)[source]¶
Run tool
namewith validatedtool_argsand return its result unmasked.This is the method a subclass implements;
call_tool()runs it and masks what it returns.ctxandtoolare keyword-only so that arguments can be added here later without breaking subclasses.