Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,9 +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

- Schema evolution now declares missing tables in place via `add_managed_table` instead of recreating the managed database. Adding a table on a later run no longer snapshots, deletes, and reloads existing data — existing tables (including dlt bookkeeping) are left untouched. Requires `hotdata-framework>=0.6.0`.
- 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.5.0] - 2026-06-29
Expand Down
5 changes: 3 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:

Expand Down
18 changes: 14 additions & 4 deletions src/hotdata_dlt_destination/merge.py
Original file line number Diff line number Diff line change
Expand Up @@ -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))}"
Expand Down
42 changes: 42 additions & 0 deletions tests/test_merge.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
Loading