From 4a9981eba8b6d9c37c8585b72ab8a19393c599c1 Mon Sep 17 00:00:00 2001 From: Eddie A Tejeda <669988+eddietejeda@users.noreply.github.com> Date: Tue, 30 Jun 2026 11:28:06 -0700 Subject: [PATCH] fix(merge): preserve unified schema on merge/upsert/insert-only combine_tables rebuilt the merged batch with pa.Table.from_pylist, which infers the schema from the first row -- silently dropping columns present only in incoming rows (existing rows sort first and lack them) and re-inferring/narrowing column types. The latter could fail the load with a 409 type conflict (e.g. decimal(12,2) -> decimal(5,2)). Build the merged batch with the unified schema of the existing and incoming tables instead. Also clarify the write-modes docs: dlt resources accept append/replace/merge only; merge is upsert-by-PK; insert-only is not selectable per resource. Verified live: a merge that adds a new column and uses a decimal column now promotes the column and preserves decimal(12,2) with no 409. --- CHANGELOG.md | 8 ++++++ README.md | 5 ++-- src/hotdata_dlt_destination/merge.py | 18 +++++++++--- tests/test_merge.py | 42 ++++++++++++++++++++++++++++ 4 files changed, 67 insertions(+), 6 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 0b5a4cc..e9c1f03 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Fixed + +- `merge`/`upsert`/`insert-only` loads no longer drop columns or narrow column types. `combine_tables` rebuilt the merged batch with `pa.Table.from_pylist`, which infers the schema from the first row — silently dropping columns present only in incoming rows (existing rows sort first and lack them) and re-inferring, often narrowing, column types. The latter could fail the load with a 409 type conflict (e.g. `decimal(12, 2)` → `decimal(5, 2)`). The merged batch is now built with the unified schema of the existing and incoming tables, so columns and types are preserved. + +### Changed + +- Clarified the write-modes documentation: dlt resources accept `append`/`replace`/`merge` only; `merge` is upsert-by-primary-key (what `upsert` resolves to), and `insert-only` is not selectable as a resource `write_disposition`. + ## [0.4.2] - 2026-06-29 diff --git a/README.md b/README.md index 4d979bd..8272e69 100644 --- a/README.md +++ b/README.md @@ -137,8 +137,9 @@ Each resource can control how its data lands in the table: |------|-------------| | `replace` | Deletes everything in the table and loads the new batch. Good for full refreshes. | | `append` | Adds new rows to the table without touching existing data. Good for event logs and immutable records. | -| `merge` (or `upsert`) | Updates existing rows by primary key, inserts new ones. Good for syncing a source of truth. | -| `insert-only` | Inserts rows whose key isn't already present; never updates existing rows. | +| `merge` (= `upsert`) | Updates existing rows by primary key, inserts new ones. Good for syncing a source of truth. | + +> dlt resources set `write_disposition` to `append`, `replace`, or `merge` only. `merge` performs upsert-by-primary-key — it is what the internal `upsert` disposition resolves to, so there is no separate `upsert` to set. The destination also implements an `insert-only` combine (insert rows whose key isn't already present, never updating existing rows), but dlt does not expose it as a resource `write_disposition`, so it cannot be selected per resource. Set the default for all resources on the destination: diff --git a/src/hotdata_dlt_destination/merge.py b/src/hotdata_dlt_destination/merge.py index 759a3bb..a68f155 100644 --- a/src/hotdata_dlt_destination/merge.py +++ b/src/hotdata_dlt_destination/merge.py @@ -73,15 +73,25 @@ def combine_tables( keys = primary_key or [fallback_key] if disposition in ("merge", "upsert"): merged = merge_rows(existing.to_pylist(), incoming.to_pylist(), primary_key=keys) - return pa.Table.from_pylist(merged) + # Build with the unified schema of both inputs. A bare ``pa.Table.from_pylist`` + # infers the schema from the first row's keys and values, which silently drops + # incoming-only columns (existing rows sort first and lack them) and re-infers + # -- often narrowing -- column types. Passing the unified schema preserves every + # column and the widest compatible type, so the merged batch never loses a + # column or narrows a type relative to the existing table. + schema = pa.unify_schemas( + [existing.schema, incoming.schema], promote_options="permissive" + ) + return pa.Table.from_pylist(merged, schema=schema) if disposition == "insert-only": existing_keys = {row_key(row, keys) for row in existing.to_pylist()} new_rows = [r for r in incoming.to_pylist() if row_key(r, keys) not in existing_keys] if not new_rows: return existing - return pa.concat_tables( - [existing, pa.Table.from_pylist(new_rows)], promote_options="permissive" - ) + # Preserve the incoming schema for the new rows (same reasoning as above), + # then let permissive concat reconcile it with the existing table. + new_table = pa.Table.from_pylist(new_rows, schema=incoming.schema) + return pa.concat_tables([existing, new_table], promote_options="permissive") raise ValueError( f"Unsupported write_disposition {disposition!r}. " f"Expected one of: {', '.join(sorted(SUPPORTED_WRITE_DISPOSITIONS))}" diff --git a/tests/test_merge.py b/tests/test_merge.py index cead65e..245f745 100644 --- a/tests/test_merge.py +++ b/tests/test_merge.py @@ -143,6 +143,48 @@ def test_combine_tables_merge_falls_back_to_dlt_id() -> None: assert result.to_pylist() == [{"_dlt_id": "a", "value": 2}] +def test_combine_tables_merge_promotes_new_column() -> None: + # Incoming introduces a column the existing table lacks. It must survive the + # merge (existing rows get null), not be dropped by first-row schema inference. + existing = _t({"id": 1, "v": "a"}, {"id": 2, "v": "b"}) + incoming = _t({"id": 2, "v": "B", "tier": "gold"}, {"id": 3, "v": "c", "tier": "silver"}) + result = combine_tables( + disposition="merge", existing=existing, incoming=incoming, primary_key=["id"] + ) + assert "tier" in result.column_names + by = {r["id"]: r for r in result.to_pylist()} + assert by[1]["tier"] is None + assert by[2]["tier"] == "gold" + assert by[3] == {"id": 3, "v": "c", "tier": "silver"} + + +def test_combine_tables_merge_preserves_existing_column_type() -> None: + # A bare from_pylist would re-infer `bal` from the values (e.g. decimal(4, 2)), + # narrowing the existing decimal(12, 2) column and breaking the load. + import decimal + + schema = pa.schema([("id", pa.int64()), ("bal", pa.decimal128(12, 2))]) + existing = pa.table({"id": [1], "bal": [decimal.Decimal("100.25")]}, schema=schema) + incoming = pa.table({"id": [1], "bal": [decimal.Decimal("50.00")]}, schema=schema) + result = combine_tables( + disposition="merge", existing=existing, incoming=incoming, primary_key=["id"] + ) + assert result.schema.field("bal").type == pa.decimal128(12, 2) + assert result.to_pylist() == [{"id": 1, "bal": decimal.Decimal("50.00")}] + + +def test_combine_tables_insert_only_promotes_new_column() -> None: + existing = _t({"id": 1, "v": "a"}) + incoming = _t({"id": 2, "v": "b", "tier": "gold"}) + result = combine_tables( + disposition="insert-only", existing=existing, incoming=incoming, primary_key=["id"] + ) + assert "tier" in result.column_names + by = {r["id"]: r for r in result.to_pylist()} + assert by[1]["tier"] is None + assert by[2]["tier"] == "gold" + + def test_combine_tables_rejects_unknown_disposition() -> None: with pytest.raises(ValueError, match="Unsupported write_disposition 'appned'"): combine_tables(