dataloader-util (0.3.4)
Installation
pip install --index-url dataloader-utilAbout this package
Config-driven CLI for loading tabular data from sources (local, S3) in formats (CSV, JSONL, Parquet) into targets (Postgres, Trino/Iceberg, Snowflake).
dataloader-util
dataloader-util is a config-driven Python CLI for moving tabular data from files or S3 into local files and analytical databases. It uses DuckDB as the in-process execution engine, so jobs can read common data formats, run SQL-style transformations, infer schemas, and write to targets from a single YAML or JSON config file.
What it does
A load job follows one pipeline:
source file(s) -> DuckDB view -> optional transforms -> target writer
A config file declares named pools of sources, targets, and transform groups. The CLI picks one of each by name at run time:
sources: { s3_csv: ..., local_parquet: ... }
targets: { warehouse: ..., archive: ... }
transforms: { clean: [...steps], camel: [...steps] }
# CLI:
dataloader-util load -c job.yaml \
--source s3_csv --target warehouse \
--transform clean --transform camel
This lets a single config file describe many pipelines (different sources, destinations, or transform chains) and pick one per run without rewriting the file.
Supported capabilities:
- Sources: local filesystem paths/globs and S3 or S3-compatible object stores.
- Input formats: CSV, JSONL/NDJSON, and Parquet.
- Targets: local/S3 Parquet or CSV files, Postgres, Trino, Trino-managed Iceberg tables, and Snowflake.
- Transforms: free-form DuckDB SQL,
WHEREfilters, column casts, and column renames. - Write modes:
replace,append, and key-basedmergefor database targets. - Config validation: strict Pydantic validation catches unknown keys and invalid combinations before a job runs.
- Lenient missing inputs: absent source files or empty globs skip successfully by default, with
fail_on_missing: trueavailable for strict pipelines. - Secrets via environment:
${VAR}placeholders are resolved from environment variables at runtime.
Requirements
- Python
>=3.14 uvfor local development commands- Docker only if you want to run local integration services/tests
- Target-specific connectivity for Postgres, Trino, Snowflake, or S3
Installation and setup
From a checkout of this repository:
uv sync
uv run dataloader-util --help
uv run dataloader-util --version
The package exposes the console command:
dataloader-util
During development, prefer running it through uv run so the repository virtual environment is used.
Quickstart: local CSV to local Parquet
Create a small input file:
mkdir -p data output
cat > data/events.csv <<'CSV'
id,event_ts,user_id,amount
1,2026-01-01T00:00:00Z,u1,12.50
2,2026-01-02T00:00:00Z,u2,-3.00
3,2026-01-03T00:00:00Z,u3,7.25
CSV
Run the included example job:
uv run dataloader-util load --config examples/job_local_to_local.yaml
With one source and one target defined, the CLI picks them
implicitly. To pick a specific source/target by name, use
--source / --target:
uv run dataloader-util load --config examples/job_multi_source_target.yaml \
--source events_csv --target events_parquet \
--transform keep_positive
Validate without running:
uv run dataloader-util validate --config examples/job_local_to_local.yaml
Perform a dry run that reads the source, applies transforms, and prints the planned target information without writing:
uv run dataloader-util load --dry-run --config examples/job_local_to_local.yaml
Enable debug logging:
uv run dataloader-util load --log-level DEBUG --config examples/job_local_to_local.yaml
CLI reference
dataloader-util [OPTIONS] COMMAND [ARGS]...
Global options:
| Option | Description |
|---|---|
--version |
Print the installed version and exit. |
--help |
Show CLI help. |
Commands:
| Command | Description |
|---|---|
load --config/-c PATH [--source NAME] [--target NAME] [--transform NAME]... |
Validate and execute a load job. |
load --dry-run --config/-c PATH [...] |
Validate, read the source, apply transforms, and report target planning without writing. |
load --log-level/-L LEVEL --config/-c PATH [...] |
Run with a Loguru level such as DEBUG, INFO, or ERROR. |
validate --config/-c PATH [...] |
Validate config shape and print a summary without reading or writing data. |
--source and --target are optional when the config defines
exactly one of each; required (and an unknown name is rejected with
a list of available names) when multiple are defined.
--transform is repeatable and optional: each occurrence names a
group from transforms: to concatenate and apply in order. Omit
--transform entirely to run with no transforms.
Config files can be YAML (.yaml, .yml) or JSON (.json).
Minimal job config
name: local_csv_to_local_parquet
sources:
events_csv:
kind: local
path: ./data/events.csv
format: csv
fail_on_missing: false # default; skip successfully if no input exists
format_options:
delimiter: ","
header: true
null_values: ["", "null"]
targets:
events_parquet:
kind: local
path: ./output/events.parquet
format: parquet
mode: replace
transforms:
keep_positive:
- kind: sql
sql: "SELECT id, event_ts, user_id, amount FROM <source> WHERE amount > 0"
Configuration reference
Top-level fields:
| Field | Required | Description |
|---|---|---|
name |
No | Optional pipeline name shown in logs and validation output. |
sources |
No | Map of named source configs. At least one is needed to run. |
targets |
No | Map of named target configs. At least one is needed to run. |
transforms |
No | Map of named transform groups; each value is an ordered list of transform steps. |
When the config has a single source, a single target, and no
transforms, the CLI runs it with no flags. With more than one of
any kind, pass --source, --target, and/or --transform to
select.
Each named target embeds its own connection: block (required for
postgres/trino/snowflake, omitted for local). Unknown keys
are rejected at every level.
Sources
Local source
sources:
events:
kind: local
path: ./data/events/*.csv
format: csv
fail_on_missing: false
format_options:
header: true
pathcan be a single file or a DuckDB-supported glob.- Exact paths and globs are checked before reading.
- If no file matches and
fail_on_missingis omitted orfalse, the run skips successfully with zero rows written. - Set
fail_on_missing: trueto fail before transforms or writes when no local input exists.
S3 source
sources:
events:
kind: s3
path: s3://my-bucket/raw/events/*.parquet
format: parquet
fail_on_missing: false
format_options:
hive_partitioning: true
region: us-east-1
endpoint_url: http://localhost:9000 # optional, for MinIO/LocalStack
access_key_id: ${AWS_ACCESS_KEY_ID} # optional
secret_access_key: ${AWS_SECRET_ACCESS_KEY} # optional
- S3 reads use DuckDB's
httpfsextension. - Exact object keys and glob patterns are checked before reading.
- Missing objects or globs with no matching objects skip successfully by default; set
fail_on_missing: trueto fail before transforms or writes. - Permission denied, bucket-not-found, credential, and network failures are always treated as errors, not as empty input.
- If explicit keys are omitted, DuckDB/AWS credential discovery is used.
- Set
endpoint_urlfor S3-compatible systems such as MinIO or LocalStack.
Formats
| Format | Reader | Options |
|---|---|---|
csv |
DuckDB read_csv_auto |
delimiter, header, null_values, compression (none, gzip, zstd) |
jsonl |
DuckDB read_json_auto with newline-delimited format |
compression (none, gzip, zstd) |
parquet |
DuckDB read_parquet |
columns, hive_partitioning |
Examples:
format: csv
format_options:
delimiter: "|"
header: true
null_values: ["", "NULL"]
compression: gzip
format: parquet
format_options:
columns: [id, event_ts, amount]
hive_partitioning: true
Transforms
A transforms: entry maps a group name to an ordered list of
transform steps. The CLI concatenates the steps from each
--transform NAME flag in the order the flags are given. Each
transform reads the previous DuckDB view and registers a new
temporary view.
SQL transform
Use <source> as the placeholder for the current input view.
transforms:
keep_positive:
- kind: sql
sql: "SELECT id, amount * 100 AS amount_cents FROM <source> WHERE amount > 0"
Filter transform
The where value is the expression only; do not include WHERE.
transforms:
active_only:
- kind: filter
where: "amount > 0 AND status = 'active'"
Rename transform
Mapping is old_name: new_name. Unmapped columns pass through unchanged.
transforms:
camel_case:
- kind: rename
mapping:
event_ts: eventTs
user_id: userId
Cast transform
Types are DuckDB SQL type names.
transforms:
cast_clean:
- kind: cast
columns:
user_id: VARCHAR
amount: DOUBLE
event_ts: TIMESTAMP
Targets
Targets live in the top-level targets: map, keyed by name. Each
target embeds its own connection: block (required for
postgres/trino/snowflake, omitted for local/s3).
Local file target
targets:
events_parquet:
kind: local
path: ./output/events.parquet
format: parquet
mode: replace
- Supported output formats:
parquet,csv. - Parent directories are created automatically.
- Use
mode: replace; local output is written with DuckDBCOPYto the configured file.
S3 target
targets:
archive_parquet:
kind: s3
path: s3://my-bucket/curated/events.parquet
format: parquet
mode: replace
region: us-east-1
endpoint_url: http://localhost:9000 # optional, for MinIO/LocalStack
access_key_id: ${AWS_ACCESS_KEY_ID} # optional
secret_access_key: ${AWS_SECRET_ACCESS_KEY} # optional
- Supported output formats:
parquet,csv. - The target format is independent from the source format, so a CSV or JSONL source can be archived as Parquet by setting
format: parqueton the target. - S3 writes use DuckDB's
httpfsextension and the same optional credential/endpoint fields as S3 sources. - If explicit keys are omitted, DuckDB/AWS credential discovery is used.
Postgres target
targets:
warehouse:
kind: postgres
table: analytics.events
mode: replace
batch_size: 50000
connection:
kind: postgres
host: localhost
port: 5432
user: etl
password: ${POSTGRES_PASSWORD}
database: analytics
schema: public
- Table names can be
tableorschema.table; a bare table defaults topublic. replacedrops and recreates the table.appendcreates the table if missing and inserts rows.mergecreates the table with a primary key if missing, stages rows in a temp table, then usesINSERT ... ON CONFLICT.- Bulk loading uses
psycopgCOPY FROM STDIN.
Trino target
targets:
events_iceberg:
kind: trino
table: iceberg.raw.events
mode: append
batch_size: 5000
connection:
kind: trino
host: localhost
port: 8080
user: etl
catalog: iceberg
schema: raw
- Table names must be
catalog.schema.table. - The target schema is created if missing.
- Writes use a temporary table plus server-side insert or merge.
- Iceberg support is provided through Trino's Iceberg connector configuration on the Trino server.
- For existing S3 Parquet files, opt into metadata-only Iceberg registration with
write_options.metadata_registration:
write_options:
metadata_registration:
enabled: true
location: s3://bucket/path/to/parquet-prefix/
format: PARQUET
recursive_directory: true
create_table_if_missing: true
table_location: s3://bucket/path/to/iceberg-table/ # optional
partitioning: [] # optional Iceberg expressions
This mode creates table metadata from the inferred source schema, then runs Trino Iceberg ALTER TABLE ... EXECUTE add_files(...) instead of inserting rows. It requires iceberg.add-files-procedure.enabled=true on the Trino Iceberg catalog, only supports s3:// Parquet locations in append or replace mode, and returns rows=0 because files are not scanned. Trino does not validate Parquet file schemas during add_files; query failures can occur later if file schemas differ from the table schema. This is for raw Parquet file registration, not system.register_table for pre-existing Iceberg metadata.
Snowflake target
targets:
events_merged:
kind: snowflake
table: RAW.EVENTS
mode: merge
merge_keys: [id]
batch_size: 100000
connection:
kind: snowflake
account: ${SNOWFLAKE_ACCOUNT}
user: ${SNOWFLAKE_USER}
password: ${SNOWFLAKE_PASSWORD}
warehouse: COMPUTE_WH
database: ANALYTICS
schema: RAW
role: ${SNOWFLAKE_ROLE}
- Table names can be
tableorschema.table; the database comes from the connection. - Writes use
snowflake-connector-pythonwrite_pandas. mergestages data into a temporary table and runs a SnowflakeMERGEstatement.- Snowflake identifiers are validated as unquoted identifier names.
Write modes
| Mode | Local | Postgres | Trino | Snowflake |
|---|---|---|---|---|
replace |
Write/overwrite output file | Drop and recreate table | Drop and recreate table | Drop table, then write with auto-create |
append |
Not intended for local files | Create table if needed, then append | Create table if needed, then append | Create table if needed, then append |
merge |
Not supported for local files | Upsert by merge_keys with ON CONFLICT |
MERGE INTO ... USING |
MERGE INTO ... USING |
For merge, target.merge_keys is required and each key must exist in the transformed source columns.
Missing source behavior
All source kinds support fail_on_missing:
sources:
daily_events:
kind: local # also works for kind: s3
path: ./data/daily/*.csv
format: csv
fail_on_missing: false
Default behavior is false: if the source exact path/object is absent or the glob matches no files/objects, the run logs source has no data; run skipped, applies no transforms, writes nothing, exits successfully, and reports zero rows written with a skipped result indicator.
Set fail_on_missing: true for strict jobs that should fail when no input exists.
Environment variable substitution
Any string value can contain ${VAR}. Placeholders are resolved before schema validation.
targets:
warehouse:
kind: postgres
table: analytics.events
connection:
kind: postgres
password: ${POSTGRES_PASSWORD}
Rules:
- Missing environment variables raise
ConfigError. - Use environment variables for credentials and account-specific values.
- Do not commit real secrets to config files.
Included examples
| File | Purpose |
|---|---|
examples/job_local_to_local.yaml |
Local CSV to local Parquet with SQL filtering. |
examples/job_local_to_local_with_transforms.yaml |
Local CSV to Parquet using filter, rename, and cast transforms. |
examples/job_local_to_postgres.yaml |
Local CSV to Postgres replace load. |
examples/job_local_to_postgres_merge.yaml |
Local CSV to Postgres upsert by key. |
examples/job_s3_to_local.yaml |
S3 Parquet to local Parquet. |
examples/job_s3_to_postgres_merge.yaml |
S3 Parquet to Postgres merge. |
examples/job_s3_to_snowflake.yaml |
S3 Parquet to Snowflake append. |
examples/job_local_to_trino_iceberg.yaml |
Local CSV to an Iceberg table through Trino. |
examples/job_s3_to_trino_iceberg_metadata.yaml |
Register existing S3 Parquet files as Iceberg metadata through Trino. |
examples/job_multi_source_target.yaml |
One file with multiple sources, targets, and transform groups; pick with --source / --target / --transform. |
Local integration services
The repository includes docker-compose.yml for local integration dependencies:
docker compose up -d
just test-integration
docker compose down
Services include:
- MinIO on
localhost:9000with console onlocalhost:9001 - Postgres on
localhost:5432(etl/etlpass/analytics) - Trino on
localhost:8080with the memory catalog
Snowflake integration tests are gated by SNOWFLAKE_* environment variables and require access to a real account.
Development
Common tasks are defined in justfile:
just # list available tasks
just dev # uv sync
just test # unit tests only
just test-integration
just test-all
just test-cov
just lint
just fmt # format check
just fmt-write # apply formatting
just types
just check # lint + types + unit tests
just fix # ruff auto-fix + format
Build artifacts can be produced with:
uv build
Troubleshooting
config file not found: check the--configpath and current working directory.unsupported config file extension: use.yaml,.yml, or.json.- Validation errors: remove unknown keys, ensure target and connection kinds match, and provide
merge_keysfor merge jobs. - Run skipped with
source has no data: the selected source path/object is absent or its glob matched nothing. This is successful by default; setfail_on_missing: trueto make it an error. - S3 read failures: verify region, credentials, bucket path, and
endpoint_urlfor S3-compatible services. Permission, bucket, credential, and network failures always error. - Trino table errors: use a fully qualified
catalog.schema.tabletarget name and verify the catalog exists on the server. - Snowflake identifier errors: use simple unquoted-style database, schema, and table names.
License
TBD