Data
Connect to SQL databases, Neo4j, SPARQL endpoints, files, and datasets through a common async interface. Work with structured schemas and query results while keeping each backend's query language.
Import connector classes, connect_data_source, and DataConnectorRegistry
from tabulaflow.data.
Connector interface
Implement DataConnector to add a custom backend.
Query failures are represented by ExecResult.error. read_only=True
requests the connector's read-only behavior; database permissions remain the
security boundary for SQL connections.
Connectors are async context managers and also expose aclose(). A
DataConnectorRegistry takes ownership of registered connectors and closes them
when its context exits.
Check ExecResult.error before using its payload. Successful statements such
as CREATE TABLE can have no DataFrame. DataFrame writes and connection
setup can raise exceptions directly.
DataConnector
Bases: Protocol
Structural interface implemented by every live queryable data source.
global_id
property
global_id: str
Return the stable source and authorization-context identity used by caches.
backend
property
backend: str
Return the concrete backend name.
read_only
property
read_only: bool
Return whether mutating operations are blocked.
run_query_async
async
run_query_async(
query: str,
parameters: Mapping[str, Any] = ...,
timeout: int | None = ...,
) -> ExecResult
Execute a query and return rows, graph data, or an error as data.
close_async
async
close_async() -> None
Permanently close the connector and release held resources.
refresh_schema_async
async
refresh_schema_async() -> DataSourceSchema
Re-introspect and return the live source schema.
Opening sources
Connection setup can raise exceptions directly.
connect_data_source
async
connect_data_source(
source: str | Sequence[str],
*,
display_name: str,
definitions: Sequence[
DataSourceDefinition
] = DEFAULT_DATA_SOURCE_DEFINITIONS,
data_dir: Path | None = None,
read_only: bool = True,
configs: DataSourceConnectorConfigs | None = None,
) -> DataConnector
Connect or load a user-facing source into a queryable connector.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
source
|
str | Sequence[str]
|
Catalog identifier, connection URL, Hugging Face dataset URL, local database path, local data-file path, or a sequence of local data-file paths. |
required |
display_name
|
str
|
Human-readable name stored in the connector schema. |
required |
definitions
|
Sequence[DataSourceDefinition]
|
Curated source definitions used for identifier and exact-locator resolution. |
DEFAULT_DATA_SOURCE_DEFINITIONS
|
data_dir
|
Path | None
|
Directory for loader-owned DuckDB files. |
None
|
read_only
|
bool
|
Whether the returned connector blocks write queries. |
True
|
configs
|
DataSourceConnectorConfigs | None
|
Backend-specific connector policies. Connector defaults are used when omitted. |
None
|
Returns:
| Type | Description |
|---|---|
DataConnector
|
A connected, queryable data connector. |
Raises:
| Type | Description |
|---|---|
FileNotFoundError
|
If a referenced local file does not exist. |
ValueError
|
If the source form is unsupported or incompatible with the other supplied sources. |
connect_url
async
connect_url(
source: str,
*,
display_name: str | None = None,
read_only: bool = True,
global_id: str | None = None,
config: SQLConnectorConfig
| Neo4jConnectorConfig
| SPARQLConnectorConfig
| None = None,
) -> DataConnector
Build the appropriate connector from an explicit connection URL.
Normalizes the URL and dispatches by explicit connector scheme. SQLAlchemy consumes SQL URL credentials inline; Neo4j and SPARQL credentials are extracted and passed separately to their drivers. SPARQL endpoints are not contacted until queried.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
source
|
str
|
A SQL, Neo4j, or explicit |
required |
display_name
|
str | None
|
Human-readable name stored in the connector schema. Inferred by the selected connector when omitted. |
None
|
read_only
|
bool
|
Request backend-appropriate read-only behavior. SQL callers still need read-only credentials or IAM for enforced security. |
True
|
global_id
|
str | None
|
Stable identity of the source and authorization context used for caching and provenance. Derived from the URL and non-secret authenticated identity, such as a username, when omitted. |
None
|
config
|
SQLConnectorConfig | Neo4jConnectorConfig | SPARQLConnectorConfig | None
|
Backend-appropriate immutable connector configuration. |
None
|
Returns:
| Type | Description |
|---|---|
DataConnector
|
A connected SQL, property-graph, or RDF connector. |
Raises:
| Type | Description |
|---|---|
ValueError
|
If the source is unsupported or required driver settings are invalid. |
TypeError
|
If |
Connector implementations
SQLConnector, Neo4jConnector, and SPARQLConnector implement the
DataConnector interface. Use their from_url_async(...) class methods to
construct them directly. SQLConnector also supports writing DataFrames to
a writable database.
SQLConnector
Schema-aware async SQL database client.
Combines:
- a :class:
ThrottledEnginefor cancellable, throttled query execution across sync and async dialects - a live :class:
SQLSchemaintrospected at construction and refreshed on demand (with optional disk cache) - read-only safety guards for borderline queries
- optional query result caching
For raw query execution without the schema / caching layer, use
:class:ThrottledEngine directly — that's what the loaders do
when building a connector from source files.
Suitable as the database layer for any tool that needs both query execution and live schema metadata: NL2SQL agents, schema browsers, query-by-example UIs, ETL jobs, catalog-aware data pipelines.
Raw SQL strings go to the DBAPI driver via exec_driver_sql,
bypassing SQLAlchemy's parameter parsing — procedural blocks
(Snowflake Scripting, PL/SQL, T-SQL batches) and dialect-specific
:identifier syntax work unchanged.
Attributes:
| Name | Type | Description |
|---|---|---|
global_id |
Stable, filename-safe identity of the source and authorization context used by caches. |
|
schema |
Current introspected SQL schema. |
|
backend |
str
|
Concrete SQLAlchemy database backend. |
language |
SQLDialect
|
SQL dialect reported by the schema. |
config |
Resolved immutable connector configuration. |
|
read_only |
Whether read-only behavior was requested. This is a client-side safety guard unless the backend enforces it natively. |
global_id
instance-attribute
global_id = global_id
schema
instance-attribute
schema = schema
config
instance-attribute
config = config
read_only
instance-attribute
read_only = read_only
backend
property
backend: str
Return the SQLAlchemy database backend name.
from_url_async
async
classmethod
from_url_async(
url: str | URL,
*,
display_name: str | None = None,
global_id: str | None = None,
read_only: bool = True,
config: SQLConnectorConfig | None = None,
schema: SQLSchema | None = None,
include_schema_names: Sequence[str] | None = None,
exclude_schema_names: Sequence[str] = (),
reuse_date_partition_schemas: bool = False,
schema_reuse_regexes: Sequence[str] = (),
sample_view_rows: bool = True,
dbms_semaphore: Semaphore | None = None,
description: str | None = None,
duckdb_init_sql: Sequence[str] = (),
connect_args: Mapping[str, Any] | None = None,
**engine_kwargs: Any,
) -> SQLConnector
Asynchronously create a SQLConnector from a database URL.
Creates a SQLAlchemy engine (async or sync) wrapped in a
:class:ThrottledEngine with concurrency control, and optionally loads
the database schema if one is not provided.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
url
|
str | URL
|
The database URL (string or :class: |
required |
display_name
|
str | None
|
Human-readable name used in
|
None
|
global_id
|
str | None
|
Globally unique, filename-safe identifier for this database connection and its caches. Derived from the URL and non-secret authenticated identity, such as its username, when omitted. In-memory databases receive a unique ID per connector instead, since they have no persistent identity. |
None
|
read_only
|
bool
|
If |
True
|
config
|
SQLConnectorConfig | None
|
Immutable connector execution and cache policy. Environment values and built-in defaults are used when omitted. |
None
|
schema
|
SQLSchema | None
|
A pre-loaded :class: |
None
|
include_schema_names
|
Sequence[str] | None
|
Optional allowlist of schemas to introspect. |
None
|
exclude_schema_names
|
Sequence[str]
|
Schemas to omit from introspection. |
()
|
reuse_date_partition_schemas
|
bool
|
If |
False
|
schema_reuse_regexes
|
Sequence[str]
|
Regexes defining additional table families whose members are asserted to share one structural schema. |
()
|
sample_view_rows
|
bool
|
Whether to read bounded row samples from views for examples, JSON inference, and previews. |
True
|
dbms_semaphore
|
Semaphore | None
|
Optional semaphore shared across connectors to limit aggregate DBMS concurrency. |
None
|
description
|
str | None
|
Optional database description stored in the schema and persisted to the schema cache. |
None
|
duckdb_init_sql
|
Sequence[str]
|
SQL statements to run on each new DuckDB connection. |
()
|
connect_args
|
Mapping[str, Any] | None
|
Options passed to each driver connection. |
None
|
**engine_kwargs
|
Any
|
Additional keyword arguments forwarded to the
SQLAlchemy engine constructor (e.g. |
{}
|
Returns:
| Type | Description |
|---|---|
SQLConnector
|
A fully initialised :class: |
SQLConnector
|
execute queries. |
release_connections_async
async
release_connections_async() -> None
Release pooled connections while keeping the connector reusable.
Raises:
| Type | Description |
|---|---|
ValueError
|
The database is in memory and releasing its connections could destroy its data. |
close_async
async
close_async() -> None
Permanently close the connector and release its resources.
For read-write DuckDB file connectors, this releases the file-level
lock so external processes (e.g. dbt run) can acquire a write
lock. Read-only file connectors already use DuckDB's native read-only
mode and do not hold a lock.
A loader-owned cleanup hook, when present, runs after the engine closes.
refresh_schema_async
async
Re-introspect the live database and update self.schema.
Use after DDL mutations (e.g. dbt run creating new tables) to
make the connector's schema reflect the current database state.
Also updates the on-disk schema cache when caching is enabled.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
tables
|
list[TableRef] | None
|
If provided, only (re-)build schemas for these tables
(or views) and merge them into the existing schema —
replacing any entry with a matching name, and appending
truly new ones. If |
None
|
Returns:
| Type | Description |
|---|---|
SQLSchema
|
The updated :class: |
write_dataframe_async
async
write_dataframe_async(
df: DataFrame,
table_name: str,
schema_name: str | None = None,
mode: TableWriteMode = "create",
) -> int
Write a DataFrame into a database table.
On success, automatically refreshes self.schema for the
target table via :meth:refresh_schema_async so the connector
reflects the new column types and table identity. DuckDB uses a
native registered relation; other dialects use pandas.to_sql.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
df
|
DataFrame
|
DataFrame to persist. |
required |
table_name
|
str
|
Destination table name. |
required |
schema_name
|
str | None
|
Optional destination schema name. |
None
|
mode
|
TableWriteMode
|
Write mode. |
'create'
|
Returns:
| Type | Description |
|---|---|
int
|
Number of rows written. |
Raises:
| Type | Description |
|---|---|
TypeError
|
If |
ValueError
|
If the connector is read-only, an argument or value is invalid, or the requested mode conflicts with table existence. |
run_query_async
async
run_query_async(
query: str | Executable,
parameters: Sequence[Any] | Mapping[str, Any] = (),
timeout: int | None | object = _UNSET,
) -> ExecResult
Execute a query and return the result.
Raw SQL strings are sent to the DBAPI driver via
exec_driver_sql, bypassing SQLAlchemy's text() parameter
parsing. This allows procedural / scripting blocks (e.g.
Snowflake Scripting DECLARE … BEGIN … END) and
:identifier patterns (e.g. VARIANT path access) to be
executed without interference.
Two surfaces stop a running query, both routed through the
dialect's cancel strategy (interrupt(), cancel(),
KILL QUERY, cursor.cancel(), …):
timeout=N— per-query deadline. On expiry the query is aborted and the returned :class:ExecResultcarrieserror.exc_type == "TimeoutError". Query-level failures (timeouts, syntax errors, …) are returned inExecResult.error, never raised.asyncio.Task.cancel()on the awaiting task — the caller wants the query to stop.CancelledErrorisBaseExceptionand propagates through this method unchanged. Wrap the call in a Task to cancel it from elsewhere (see Example below).
See :meth:ThrottledEngine.execute_async for the full
transaction / rollback contract, including the caveat that
MySQL and Oracle DDL auto-commit per statement — cancel
stops execution but cannot undo committed effects on those
dialects.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
query
|
str | Executable
|
A raw SQL string or a SQLAlchemy |
required |
parameters
|
Sequence[Any] | Mapping[str, Any]
|
Bind parameters. For raw SQL strings these must
use the driver's native paramstyle (e.g. |
()
|
timeout
|
int | None | object
|
Query timeout in seconds. When omitted, use the connector
configuration; |
_UNSET
|
Returns:
| Name | Type | Description |
|---|---|---|
An |
ExecResult
|
class: |
ExecResult
|
an error) and latency information. |
Raises:
| Type | Description |
|---|---|
CancelledError
|
If the awaiting task was cancelled.
Propagates as-is; not wrapped in :class: |
Example
Cancel a long-running query from elsewhere (e.g. a Ctrl+C handler or an external trigger):
.. code-block:: python
task = asyncio.create_task(
connector.run_query_async("SELECT ... long-running")
)
# ... on Ctrl+C / external trigger:
task.cancel()
try:
result = await task
except asyncio.CancelledError:
... # query was aborted server-side
TableWriteMode
module-attribute
TableWriteMode: TypeAlias = Literal[
"create", "append", "replace_rows", "replace_table"
]
Neo4jConnector
Property-graph connector for Neo4j databases.
Uses the official neo4j async Python driver (Bolt protocol).
Fast schema introspection uses Neo4j metadata procedures; full-scan mode
exhaustively derives observed properties and topology from graph data.
Attributes:
| Name | Type | Description |
|---|---|---|
global_id |
Stable, filename-safe identity of the source and authorization context used by caches. |
|
schema |
Current introspected property-graph schema. |
|
backend |
Literal['neo4j']
|
Graph database backend name ( |
language |
GraphQueryLanguage
|
Graph query language ( |
config |
Resolved immutable connector configuration. |
|
read_only |
Whether sessions use server-enforced read access. |
backend
class-attribute
backend: Literal['neo4j'] = 'neo4j'
global_id
instance-attribute
global_id = global_id
schema
instance-attribute
schema = schema
config
instance-attribute
config = config
read_only
instance-attribute
read_only = read_only
from_url_async
async
classmethod
from_url_async(
url: str,
*,
global_id: str | None = None,
auth: tuple[str, str] | Auth | None = None,
database: str | None = None,
display_name: str | None = None,
schema: PropertyGraphSchema | None = None,
read_only: bool = True,
config: Neo4jConnectorConfig | None = None,
**driver_kwargs: Any,
) -> Neo4jConnector
Create a connector from a Neo4j Bolt URL.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
url
|
str
|
Neo4j URL, such as |
required |
global_id
|
str | None
|
Globally unique, filename-safe identifier for this database connection and its caches. Derived from the URL, authenticated identity, and database when omitted. |
None
|
auth
|
tuple[str, str] | Auth | None
|
|
None
|
database
|
str | None
|
Neo4j database name. |
None
|
display_name
|
str | None
|
Human-readable name used in
|
None
|
schema
|
PropertyGraphSchema | None
|
Pre-loaded schema. If |
None
|
read_only
|
bool
|
Use Neo4j's server-enforced read access mode when
|
True
|
config
|
Neo4jConnectorConfig | None
|
Immutable connector execution and cache policy. Environment values and built-in defaults are used when omitted. |
None
|
**driver_kwargs
|
Any
|
Extra keyword arguments for
|
{}
|
run_query_async
async
run_query_async(
query: str,
parameters: Mapping[str, Any] | None = None,
timeout: int | None | object = _UNSET,
) -> ExecResult
Execute Cypher and return tabular data with an optional graph.
Query failures are returned in ExecResult.error. Omitting
timeout uses the connector configuration; None disables it.
Task cancellation propagates.
close_async
async
close_async() -> None
Close the Neo4j driver and all pooled connections.
refresh_schema_async
async
refresh_schema_async() -> PropertyGraphSchema
Re-introspect the live database, bypassing cache on read.
SPARQLConnector
Read-only connector for a standards-compatible SPARQL query endpoint.
backend
class-attribute
backend: Literal['sparql'] = 'sparql'
endpoint_url
instance-attribute
endpoint_url = endpoint_url
global_id
instance-attribute
global_id = global_id
schema
instance-attribute
schema = schema
read_only
instance-attribute
read_only = True
config
instance-attribute
config = config
from_url_async
async
classmethod
from_url_async(
url: str,
*,
display_name: str | None = None,
global_id: str | None = None,
read_only: bool = True,
auth: tuple[str, str] | None = None,
config: SPARQLConnectorConfig | None = None,
description: str | None = None,
transport: AsyncBaseTransport | None = None,
) -> SPARQLConnector
Create a connector for an HTTP SPARQL query endpoint.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
url
|
str
|
Absolute HTTP or HTTPS query-endpoint URL. |
required |
display_name
|
str | None
|
Human-readable name stored in the RDF schema. Defaults to the endpoint hostname. |
None
|
global_id
|
str | None
|
Stable cache and source identifier derived from |
None
|
read_only
|
bool
|
Must remain true until SPARQL Update is supported. |
True
|
auth
|
tuple[str, str] | None
|
Optional HTTP Basic username and password. |
None
|
config
|
SPARQLConnectorConfig | None
|
Immutable execution and HTTP policy. |
None
|
description
|
str | None
|
Optional source description stored in the RDF schema. |
None
|
transport
|
AsyncBaseTransport | None
|
Optional HTTPX transport, primarily for custom networking and deterministic tests. |
None
|
Returns:
| Type | Description |
|---|---|
SPARQLConnector
|
A SPARQL connector without contacting the endpoint. |
Raises:
| Type | Description |
|---|---|
ValueError
|
If the URL or requested access mode is unsupported. |
run_query_async
async
run_query_async(
query: str,
parameters: Mapping[str, Any] | None = None,
timeout: int | None | object = _UNSET,
) -> ExecResult
Execute a SPARQL SELECT or ASK query and return a tabular result.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
query
|
str
|
Complete SPARQL query text. |
required |
parameters
|
Mapping[str, Any] | None
|
Must be empty because SPARQL has no standard parameter binding protocol. |
None
|
timeout
|
int | None | object
|
Overall operation timeout. Omitting it uses the configured
default; |
_UNSET
|
Returns:
| Type | Description |
|---|---|
ExecResult
|
Query rows or failure details as an execution result. Task |
ExecResult
|
cancellation propagates. |
refresh_schema_async
async
refresh_schema_async() -> RDFSchema
Return the configured minimal RDF schema without network discovery.
close_async
async
close_async() -> None
Close the endpoint HTTP client.
Configuration
Connector settings use explicit arguments first, then TABULAFLOW_*
environment variables, then defaults. Pass DataSourceConnectorConfigs to
connect_data_source when the source's backend is selected at runtime.
SQL connectors support both sync and async drivers through the same awaited
API. A query's timeout argument overrides the configured deadline and
requests cancellation of the underlying query. Timeouts appear in
ExecResult.error, while task cancellation propagates as
asyncio.CancelledError.
Schema and query caches are off by default. Enable them for reusable database
snapshots; call refresh_schema_async() to refresh the schema explicitly.
Query caching requires a read-only connector. Share a dbms_semaphore across
SQL connectors to limit concurrency against the same warehouse.
DataSourceConnectorConfigs
dataclass
DataSourceConnectorConfigs(
sql: SQLConnectorConfig = SQLConnectorConfig(),
neo4j: Neo4jConnectorConfig = Neo4jConnectorConfig(),
sparql: SPARQLConnectorConfig = SPARQLConnectorConfig(),
)
Backend-specific policies used when connecting an untyped data source.
sql
class-attribute
instance-attribute
sql: SQLConnectorConfig = field(
default_factory=SQLConnectorConfig
)
neo4j
class-attribute
instance-attribute
neo4j: Neo4jConnectorConfig = field(
default_factory=Neo4jConnectorConfig
)
sparql
class-attribute
instance-attribute
sparql: SPARQLConnectorConfig = field(
default_factory=SPARQLConnectorConfig
)
SQLConnectorConfig
Bases: _CachedSchemaConnectorConfig
Operational policy for a SQL connector.
Attributes:
| Name | Type | Description |
|---|---|---|
cache_dir |
Path
|
Root directory for schema and query-result caches. |
max_result_rows |
PositiveInt | None
|
Maximum rows materialized by one query, or |
query_timeout_seconds |
PositiveInt | None
|
Default query timeout, or |
max_query_concurrency |
PositiveInt
|
Maximum in-flight queries and connection-pool size per connector. |
schema_cache_mode |
Literal['off', 'read_write', 'refresh', 'cache_only']
|
Schema cache read/write policy. |
sql_column_stats_enabled |
bool
|
Whether to collect exact row counts and column statistics for physical tables. Tables and views are always enriched from one bounded row sample; views are never exhaustively profiled. |
sql_query_cache_mode |
Literal['off', 'read_write', 'refresh']
|
Query-result cache read/write policy. |
sql_column_stats_enabled
class-attribute
instance-attribute
sql_column_stats_enabled: bool = False
sql_query_cache_mode
class-attribute
instance-attribute
sql_query_cache_mode: Literal[
"off", "read_write", "refresh"
] = "off"
max_result_rows
class-attribute
instance-attribute
max_result_rows: PositiveInt | None = 1000000
query_timeout_seconds
class-attribute
instance-attribute
query_timeout_seconds: PositiveInt | None = 300
max_query_concurrency
class-attribute
instance-attribute
max_query_concurrency: PositiveInt = 8
cache_dir
class-attribute
instance-attribute
cache_dir: Path = DEFAULT_CACHE_DIR
schema_cache_mode
class-attribute
instance-attribute
schema_cache_mode: Literal[
"off", "read_write", "refresh", "cache_only"
] = "off"
Neo4jConnectorConfig
Bases: _CachedSchemaConnectorConfig
Operational policy for a Neo4j connector.
Attributes:
| Name | Type | Description |
|---|---|---|
cache_dir |
Path
|
Root directory for schema caches. |
max_result_rows |
PositiveInt | None
|
Maximum rows materialized by one query, or |
query_timeout_seconds |
PositiveInt | None
|
Default query timeout, or |
max_query_concurrency |
PositiveInt
|
Maximum in-flight queries and connection-pool size per connector. |
schema_cache_mode |
Literal['off', 'read_write', 'refresh', 'cache_only']
|
Schema cache read/write policy. |
graph_schema_introspection_mode |
Literal['fast', 'full_scan']
|
|
max_graph_result_nodes |
PositiveInt | None
|
Maximum nodes extracted into a graph result. |
max_graph_result_edges |
PositiveInt | None
|
Maximum edges extracted into a graph result. |
graph_schema_introspection_mode
class-attribute
instance-attribute
graph_schema_introspection_mode: Literal[
"fast", "full_scan"
] = "fast"
max_graph_result_nodes
class-attribute
instance-attribute
max_graph_result_nodes: PositiveInt | None = 300
max_graph_result_edges
class-attribute
instance-attribute
max_graph_result_edges: PositiveInt | None = 700
max_result_rows
class-attribute
instance-attribute
max_result_rows: PositiveInt | None = 1000000
query_timeout_seconds
class-attribute
instance-attribute
query_timeout_seconds: PositiveInt | None = 300
max_query_concurrency
class-attribute
instance-attribute
max_query_concurrency: PositiveInt = 8
cache_dir
class-attribute
instance-attribute
cache_dir: Path = DEFAULT_CACHE_DIR
schema_cache_mode
class-attribute
instance-attribute
schema_cache_mode: Literal[
"off", "read_write", "refresh", "cache_only"
] = "off"
SPARQLConnectorConfig
Bases: _ConnectorConfig
Operational policy for a SPARQL connector.
Attributes:
| Name | Type | Description |
|---|---|---|
max_result_rows |
PositiveInt | None
|
Maximum rows materialized by one query, or |
query_timeout_seconds |
PositiveInt | None
|
Overall query deadline, including throttling and
retries, or |
max_query_concurrency |
PositiveInt
|
Maximum in-flight queries per connector. |
max_sparql_response_bytes |
PositiveInt
|
Maximum decompressed SPARQL response bytes buffered. |
max_sparql_response_bytes
class-attribute
instance-attribute
max_sparql_response_bytes: PositiveInt = 50 * 1024 * 1024
max_result_rows
class-attribute
instance-attribute
max_result_rows: PositiveInt | None = 1000000
query_timeout_seconds
class-attribute
instance-attribute
query_timeout_seconds: PositiveInt | None = 300
max_query_concurrency
class-attribute
instance-attribute
max_query_concurrency: PositiveInt = 8
Connector registry
Registry validation can raise exceptions directly.
DataConnectorRegistry
DataConnectorRegistry()
Own named data connectors for a runtime.
has
has(alias: str) -> bool
Return whether a connector is registered for alias.
get
get(alias: str) -> DataConnector
Return the connector registered for alias.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
alias
|
str
|
The connector alias. |
required |
Raises:
| Type | Description |
|---|---|
ValueError
|
If |
register
register(alias: str, connector: DataConnector) -> None
Register a connector under alias and take ownership of it.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
alias
|
str
|
The connector alias. |
required |
connector
|
DataConnector
|
The connector instance to store. |
required |
Raises:
| Type | Description |
|---|---|
ValueError
|
If |
close_async
async
close_async(alias: str) -> bool
Remove and close the connector registered for alias.
list_aliases
list_aliases() -> list[str]
Return the registered connector aliases.
close_all_async
async
close_all_async() -> None
Remove and close every connector, reporting all close failures.
validate_global_id
validate_global_id(global_id: str) -> str
Validate and return a globally unique, filename-safe connector ID.
File and dataset loaders
connect_data_source dispatches to these loaders. Call them directly when you
need loader-specific options.
load_files
async
load_files(
global_id: str,
file_paths: list[str],
*,
display_name: str | None = None,
data_dir: str | None = None,
read_only: bool = True,
config: SQLConnectorConfig | None = None,
) -> SQLConnector
Create a connector from CSV, Excel, Parquet, or JSON files.
Each file is loaded into a DuckDB table. Supported formats:
.csv, .tsv, .xlsx, .xls, .parquet, .json,
.jsonl, .ndjson.
File ownership: load_files claims the path
<data_dir>/<global_id>.duckdb (or a fresh temp file when
data_dir is None). Any existing file at that path is silently
deleted before the load — the DB is always (re)built from
scratch. Don't point data_dir+global_id at a file you want to
preserve. When data_dir is None the temp file is also deleted
by :meth:SQLConnector.close_async.
Concurrency precondition: the caller must ensure (data_dir,
global_id) is unique across live :class:SQLConnector instances in
the process. Two simultaneous load_files calls resolving to the
same path will corrupt each other (the second's unlink-existing
step deletes the first's open DB file). The tabulaflow CLI guarantees
this via the alias-uniqueness check on /connect.
Atomicity: load_files either returns a successfully-loaded
:class:SQLConnector or leaves no DB file at the target path. Any
failure mid-load (cancellation, exception, error from the loader)
unlinks the partial file before re-raising. Cancellation
propagates through the SQLConnector's standard query path —
DuckDB's conn.interrupt() aborts the in-flight CREATE TABLE.
Read-only enforcement: the DuckDB connection is always opened
read-write (DDL is required for the load). read_only=True is
enforced at the SQLConnector layer — write statements via
:meth:SQLConnector.run_query_async are blocked. This is sufficient
when the calling session is the sole owner of the cache file.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
global_id
|
str
|
Globally unique identifier for this connection, also used as the cache key when loading the schema. |
required |
file_paths
|
list[str]
|
Paths to data files to load. |
required |
display_name
|
str | None
|
Display name for the data source. Defaults to the first file's stem. |
None
|
data_dir
|
str | None
|
Directory to store the DuckDB file. If |
None
|
read_only
|
bool
|
If True, block write statements at the SQLConnector layer. The underlying DuckDB connection is always opened read-write so the loader can issue CREATE TABLE statements. |
True
|
config
|
SQLConnectorConfig | None
|
Immutable connector execution and cache policy. Environment values and built-in defaults are used when omitted. |
None
|
Returns:
| Name | Type | Description |
|---|---|---|
A |
SQLConnector
|
class: |
load_hf_dataset
async
load_hf_dataset(
dataset_url: str,
*,
display_name: str | None = None,
read_only: bool = True,
config: SQLConnectorConfig | None = None,
) -> SQLConnector
Load a HuggingFace dataset into a DuckDB-backed SQLConnector.
Small datasets are fully materialized. Large datasets get a lazy view
over all parquet files plus a materialized sample table.
DuckDB files are cached in ~/.tabulaflow/cache/hf/ across sessions.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
dataset_url
|
str
|
A HuggingFace dataset URL. |
required |
display_name
|
str | None
|
Display name for the data source. Defaults to the dataset name. |
None
|
read_only
|
bool
|
If True, block write statements. |
True
|
config
|
SQLConnectorConfig | None
|
Immutable connector execution and cache policy. Environment values and built-in defaults are used when omitted. |
None
|
Returns:
| Name | Type | Description |
|---|---|---|
A |
SQLConnector
|
class: |
HuggingFaceSubsetRequiredError
HuggingFaceSubsetRequiredError(
dataset_id: str, subsets: Sequence[str]
)
Bases: ValueError
A dataset has multiple subsets and none was specified.
dataset_id
instance-attribute
dataset_id = dataset_id
subsets
instance-attribute
subsets = tuple(subsets)
DATA_FILE_EXTENSIONS
module-attribute
DATA_FILE_EXTENSIONS = frozenset(
{
".csv",
".tsv",
".xlsx",
".xls",
".parquet",
".json",
".jsonl",
".ndjson",
}
)
File suffixes accepted by :func:load_files.
build_hf_dataset_url
build_hf_dataset_url(
dataset_id: str, subset: str, split: str | None = None
) -> str
Build a canonical Hugging Face viewer URL.
parse_hf_dataset_url
parse_hf_dataset_url(
url: str,
) -> tuple[str, str | None, str | None]
Parse a HuggingFace dataset URL into (dataset_id, subset, split).
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
url
|
str
|
A URL like |
required |
Returns:
| Type | Description |
|---|---|
str
|
A tuple of |
str | None
|
may be |
Raises:
| Type | Description |
|---|---|
ValueError
|
If the URL does not match the expected HuggingFace dataset pattern. |
is_hf_dataset_url
is_hf_dataset_url(url: str) -> bool
Return True if the URL points to a HuggingFace dataset.
Source catalog
Catalog entries are named source definitions, not live connectors.
connect_data_source resolves an entry to its source, then dispatches to
the appropriate loader or connector. An entry can name a local file, a
Hugging Face dataset, or a connection URL. The connector registry instead
holds live connector instances.
DataSourceDefinition
dataclass
DataSourceDefinition(
id: str, source: str, description: str
)
A curated source available by a stable catalog identifier.
id
instance-attribute
id: str
source
instance-attribute
source: str
description
instance-attribute
description: str
DEFAULT_DATA_SOURCE_DEFINITIONS
module-attribute
DEFAULT_DATA_SOURCE_DEFINITIONS = (WIKIDATA,)
resolve_data_source_definition
resolve_data_source_definition(
source: str,
definitions: Sequence[
DataSourceDefinition
] = DEFAULT_DATA_SOURCE_DEFINITIONS,
) -> DataSourceDefinition | None
Resolve a curated source by identifier or exact normalized locator.
Execution errors
These errors identify result-size and SPARQL response failures. Query methods
that return ExecResult capture execution failures in error; lower-level
execution APIs may raise directly.
ResultTooLargeError
ResultTooLargeError(max_rows: int)
Bases: RuntimeError
A query produced more rows than may be materialized safely.
InvalidSPARQLResultError
Bases: ValueError
A response is not a valid supported SPARQL JSON result.
SPARQLResponseTooLargeError
SPARQLResponseTooLargeError(max_bytes: int)
Bases: RuntimeError
A response exceeded the configured byte limit.