Skip to content

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.

schema property

Return the connector's current source schema.

backend property

backend: str

Return the concrete backend name.

language property

language: QueryLanguage

Return the query language understood by the connector.

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 sparql+http(s) connection URL. Examples include postgresql://user:pass@host/db, bigquery://project/dataset, neo4j+s://user:pass@host, sparql+https://query.wikidata.org/sparql, sqlite+aiosqlite:///data.sqlite, and duckdb:///data.duckdb.

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 config does not match the URL backend.

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:ThrottledEngine for cancellable, throttled query execution across sync and async dialects
  • a live :class:SQLSchema introspected 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.

language property

language: SQLDialect

Return the connector's SQL dialect.

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:SQLAlchemyURL).

required
display_name str | None

Human-readable name used in schema.display_name. Defaults to the supplied schema's name, then the URL's database value, then the backend name.

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 (the default), write statements (INSERT, UPDATE, DELETE, DROP, etc.) recognized by the client guard are rejected before reaching the database. This is not a security boundary; use read-only credentials or IAM for enforcement.

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:SQLSchema. When None the schema is loaded (and cached) automatically via :func:_load_schema_async.

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 True, date-suffixed table families reuse one representative's structural schema.

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. pool_pre_ping).

{}

Returns:

Type Description
SQLConnector

A fully initialised :class:SQLConnector instance ready to

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

refresh_schema_async(
    tables: list[TableRef] | None = None,
) -> SQLSchema

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, do a full rebuild.

None

Returns:

Type Description
SQLSchema

The updated :class:SQLSchema.

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 creates a new table, append adds rows, replace_rows replaces rows while preserving the table definition, and replace_table recreates the table.

'create'

Returns:

Type Description
int

Number of rows written.

Raises:

Type Description
TypeError

If df is not a pandas DataFrame.

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:ExecResult carries error.exc_type == "TimeoutError". Query-level failures (timeouts, syntax errors, …) are returned in ExecResult.error, never raised.
  • asyncio.Task.cancel() on the awaiting task — the caller wants the query to stop. CancelledError is BaseException and 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 Executable.

required
parameters Sequence[Any] | Mapping[str, Any]

Bind parameters. For raw SQL strings these must use the driver's native paramstyle (e.g. %(name)s for pyformat drivers).

()
timeout int | None | object

Query timeout in seconds. When omitted, use the connector configuration; None explicitly disables the timeout. On expiry, the result's error.exc_type is "TimeoutError" (not raised).

_UNSET

Returns:

Name Type Description
An ExecResult

class:ExecResult containing the result DataFrame (or

ExecResult

an error) and latency information.

Raises:

Type Description
CancelledError

If the awaiting task was cancelled. Propagates as-is; not wrapped in :class:ExecResult.

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 (neo4j).

language GraphQueryLanguage

Graph query language (cypher).

config

Resolved immutable connector configuration.

read_only

Whether sessions use server-enforced read access.

backend class-attribute

backend: Literal['neo4j'] = 'neo4j'

language class-attribute

language: GraphQueryLanguage = 'cypher'

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 "bolt://localhost:7687" for a local direct connection or "neo4j+s://host" for hosted routing with trusted TLS. Preserve the deployment-provided scheme.

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

(username, password) tuple or neo4j.Auth object.

None
database str | None

Neo4j database name. None uses the server default.

None
display_name str | None

Human-readable name used in schema.display_name. Defaults to the supplied schema's name, then the requested or server-default database name, then "Neo4j".

None
schema PropertyGraphSchema | None

Pre-loaded schema. If None, the schema is introspected automatically.

None
read_only bool

Use Neo4j's server-enforced read access mode when True.

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 neo4j.AsyncGraphDatabase.driver.

{}

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'

language class-attribute

language: GraphQueryLanguage = '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 url and Basic Auth username when omitted.

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; None disables it.

_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

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 None for no limit.

query_timeout_seconds PositiveInt | None

Default query timeout, or None to disable.

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 None for no limit.

query_timeout_seconds PositiveInt | None

Default query timeout, or None to disable.

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']

fast for metadata procedures or full_scan for observed graph data.

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 None for no limit.

query_timeout_seconds PositiveInt | None

Overall query deadline, including throttling and retries, or None to disable.

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 alias is not registered.

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 alias is already registered.

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, a system temp directory is used and the file is unlinked on close.

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:SQLConnector backed by a DuckDB database.

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:SQLConnector backed by a DuckDB database.

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 https://huggingface.co/datasets/user/name or https://huggingface.co/datasets/user/name/viewer/subset/split.

required

Returns:

Type Description
str

A tuple of (dataset_id, subset, split) where subset and split

str | None

may be None.

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.