ETL & Data Pipelines#
SQLSpec works well for ETL (Extract-Transform-Load) workflows. Use multiple database configs to move data between systems, and leverage Arrow-based methods for high-performance bulk transfers.
Multi-Database ETL#
Register source and target databases on a single SQLSpec instance. Extract
from one, transform in Python, and load into the other.
multi-database ETL pipeline#from sqlspec import SQLSpec
from sqlspec.adapters.sqlite import SqliteConfig
# Simulate an ETL pipeline with two SQLite databases
spec = SQLSpec()
source_config = spec.add_config(SqliteConfig(connection_config={"database": str(tmp_path / "source.db")}))
target_config = spec.add_config(SqliteConfig(connection_config={"database": str(tmp_path / "target.db")}))
# Step 1: Seed source data
with spec.provide_session(source_config) as session:
session.execute("create table orders (id integer primary key, amount real, status text)")
session.execute_many(
"insert into orders (amount, status) values (?, ?)",
[(100.0, "complete"), (50.0, "pending"), (200.0, "complete")],
)
# Step 2: Extract from source
with spec.provide_session(source_config) as session:
completed = session.select("select id, amount from orders where status = ?", "complete")
# Step 3: Load into target
with spec.provide_session(target_config) as session:
session.execute("create table revenue (order_id integer, amount real)")
session.execute_many(
"insert into revenue (order_id, amount) values (?, ?)", [(row["id"], row["amount"]) for row in completed]
)
total = session.select_value("select sum(amount) from revenue")
print(f"Total revenue: {total}") # Total revenue: 300.0
Arrow-Based Bulk Transfer#
For large datasets, use select_to_arrow() to get results as Apache Arrow
tables. This avoids per-row Python object overhead and enables zero-copy
transfers between databases that support native Arrow (ADBC, DuckDB, BigQuery,
Arrow ODBC, mssql-python, and Oracle).
Passing native_only=True ensures that the query executes through the driver's
fast, zero-copy C/C++ columnar path. If the adapter does not support native Arrow,
it raises ImproperConfigurationError rather than
silently falling back to row-by-row dict conversion.
# Extract as Arrow table using zero-copy native path
arrow_result = await source_session.select_to_arrow(
"SELECT * FROM large_table WHERE updated > :since",
{"since": last_sync},
native_only=True,
)
# Convert to Polars or Pandas DataFrames
df_polars = arrow_result.to_polars()
df_pandas = arrow_result.to_pandas()
# Direct export to cloud storage (Parquet / Arrow IPC)
await arrow_result.write_to_storage_async("s3://my-bucket/exports/table.parquet")
Direct Zero-Copy Cross-Database Transfer#
Because load_from_arrow() accepts an ArrowResult directly, you can stream
data between heterogeneous databases without any intermediate Python objects:
# 1. Zero-copy extract from PostgreSQL via ADBC
arrow_result = await postgres_session.select_to_arrow(
"SELECT id, event_type, payload, created_at FROM events WHERE processed = false",
native_only=True,
)
# 2. Direct zero-copy load into DuckDB analytical staging table
job = duckdb_session.load_from_arrow("staging_events", arrow_result)
print(f"Transferred {job.telemetry['rows_processed']} rows")
Streaming Large Result Sets#
To stream datasets that exceed available RAM, use return_format="reader" to
obtain a pyarrow.RecordBatchReader:
arrow_result = await session.select_to_arrow(
"SELECT * FROM events",
return_format="reader", # Streams RecordBatch objects
batch_size=10000, # Rows per batch
)
reader = arrow_result.get_data()
for batch in reader:
process_batch(batch)
Supported return_format values:
"table"-- singlepyarrow.Table(default)"batch"-- singleRecordBatch"batches"-- iterator ofRecordBatchobjects"reader"--RecordBatchReaderfor streaming
See also
Bulk Ingest for the inbound side -- loading Arrow tables, staged files, and in-memory records into a table via native driver primitives.
DuckDB as Staging Layer#
DuckDB excels as an ETL staging layer because it can read Parquet, CSV, and JSON files natively and attach to external PostgreSQL databases.
from sqlspec import SQLSpec
from sqlspec.adapters.duckdb import DuckDBConfig
spec = SQLSpec()
config = spec.add_config(
DuckDBConfig(connection_config={"database": "/tmp/staging.db"})
)
with spec.provide_session(config) as session:
# Read directly from Parquet files
session.execute(
"CREATE TABLE staging AS SELECT * FROM read_parquet('data/*.parquet')"
)
# Transform and aggregate
session.execute(
"CREATE TABLE summary AS "
"SELECT date, COUNT(*) as events "
"FROM staging GROUP BY date"
)
Direct Export to Cloud Storage#
For massive analytical tables, export directly to Parquet or CSV on AWS S3, Google Cloud Storage, or Azure Blob Storage without loading rows into Python memory:
with spec.provide_session(config) as session:
session.select_to_storage(
"SELECT * FROM sales_fact WHERE transaction_date >= :cutoff",
destination="s3://data-lake/exports/sales_2026.parquet",
file_format="parquet",
cutoff="2026-01-01",
)