Add Parquet and DuckDB output formats - #333
Conversation
Transformed objects stream to a staging JSONL file, then DuckDB's read_json_auto builds the artifact. That inference is the whole point: a nested object becomes a STRUCT column and a list of them STRUCT[], so consumers read a nested field with a column access instead of a join, without linkml-map describing the shape. Staging keeps the streaming contract — objects are never all held in memory, and the output path only ever contains the finished artifact. `map-data -o out.parquet` and `-o study.duckdb` work by extension, and both are usable as -O targets. DuckDB output writes one table per run, named by --table-name or the file stem, so repeated runs build up a database covering several target classes. StreamWriter gains staging_path/finalize_path hooks. The tabular header rewrite moves onto finalize_path, so MultiStreamWriter no longer type-checks for TabularStreamWriter to decide on post-processing. Closes #325
There was a problem hiding this comment.
Pull request overview
This PR adds first-class columnar outputs (Parquet and DuckDB) to linkml-map’s streaming writer pipeline by staging transformed objects as JSONL and letting DuckDB infer nested types (STRUCT / STRUCT[]), enabling nested-field access without joins and without holding all rows in memory.
Changes:
- Introduces
ParquetStreamWriterandDuckDBStreamWriterbuilt on a newColumnarStreamWriterstaging + finalize flow. - Refactors
StreamWriter/MultiStreamWriterto support per-writer staging paths andfinalize_path()post-processing (including moving tabular header rewrites there). - Adds tests covering nesting preservation, staging cleanup, DuckDB table naming, and multi-output fan-out.
Reviewed changes
Copilot reviewed 3 out of 3 changed files in this pull request and generated 2 comments.
| File | Description |
|---|---|
| tests/test_writers/test_columnar_output.py | Adds coverage for Parquet/DuckDB outputs, nesting types, staging cleanup, and multi-run/table behavior. |
| src/linkml_map/writers/output_streams.py | Adds new output formats and columnar writers; introduces staging/finalize hooks and refactors post-processing. |
| src/linkml_map/cli/cli.py | Threads --table-name through streaming CLI path and uses updated writer creation for primary output. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
Apply --table-name to every DuckDB writer a command builds, not just the primary output: -O out.duckdb --table-name X was silently falling back to the file stem. Route parquet/duckdb through the writer path on single-object input too. Format inference fires on the output extension regardless of input shape, so `-o out.parquet data.yaml` reached dump_output and raised NotImplementedError. These artifacts are binary files DuckDB writes, so a missing -o now fails as a CLI error rather than attempting stdout. Add both formats to --output-format's choices, which otherwise accepted them by extension inference but rejected them when named explicitly.
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 4 out of 4 changed files in this pull request and generated no new comments.
Suppressed comments (2)
src/linkml_map/writers/output_streams.py:642
ColumnarStreamWriter.finalize_path()only unlinks the staging JSONL on success. If DuckDB conversion raises, the.staging.jsonlfile is left behind (potentially large/sensitive) and the docstring’s “then remove the staging file” guarantee no longer holds. Consider cleaning up the staging file in afinallyblock (logging a warning if removal fails) so failures don’t strand temp artifacts.
try:
connection.execute(self._conversion_sql(target, staged))
finally:
connection.close()
staged.unlink()
src/linkml_map/writers/output_streams.py:639
- DuckDB conversion here uses a bare
duckdb.connect(...), but other parts of the codebase intentionally uselinkml_map.utils.lookup_index.make_connection()to apply container/cgroup-aware settings (seesrc/linkml_map/transformer/engine.py:119andsrc/linkml_map/utils/lookup_index.py:160-166). For large JSON→Parquet/DB conversions, skipping those settings can increase the risk of OOM or unpredictable resource usage in constrained environments. Consider reusing the same connection configuration for these conversions (including the file-backed case).
This issue also appears on line 638 of the same file.
import duckdb
connection = duckdb.connect(*self._connect_args(target))
try:
connection.execute(self._conversion_sql(target, staged))
Closes #325.
Transformed objects stream to a staging JSONL file, which DuckDB's
read_json_autothen turns into the artifact. That inference is the point — it types each field on its own, so a nested object lands as aSTRUCTcolumn and a list of them asSTRUCT[]:A consumer reads a nested field with a column access rather than a join, and linkml-map never has to describe the shape. Staging preserves the streaming contract: objects are never all held in memory, and the output path only ever contains the finished artifact.
Both formats work by extension and as
-Otargets alongside a text primary output. DuckDB output writes one table per run — named by--table-name, defaulting to the file stem — so repeated runs against one path build a database covering several target classes, which is the shape downstream consumers assemble by hand today.duckdbwas already a hard dependency (the join engine uses it), so this adds none.Why this rather than the SQL backend
#325 framed Parquet as work on
DuckDBTransformer. It isn't:SQLCompilersupports almost none of the specification language —expr,value,case(), unit conversion and value mappings are all dropped silently — so nothing real transforms through it. The Python backend already produces correct nested output; only the serialization step was missing. See the closing comment on #332 for the full finding.Incidental refactor
StreamWritergains two optional hooks,staging_pathandfinalize_path. The tabular header rewrite moves ontofinalize_path, soMultiStreamWriterno longer type-checks forTabularStreamWriterto decide whether to post-process.