Skip to content
Closed
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
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,17 @@ def is_transient_delta_maintenance_error(error: BaseException) -> bool:
)


# delta-rs runs a post-commit hook after every write/merge that, per table defaults
# (`delta.enableExpiredLogCleanup=true`, 30-day `logRetentionDuration`), batch-deletes
# `_delta_log` JSON files once a table's commit history crosses the retention window — via the
# same object_store bulk DeleteObjects call that TRANSIENT_OBJECT_STORE_ERRORS guards elsewhere.
# Our storage backend can return a DeleteObjects response delta-rs's XML parser doesn't
# recognize ("unknown variant `Code`, expected `Deleted` or `Error`"), which fails the whole
# commit even though the write itself already succeeded. We don't rely on this automatic
# cleanup (nothing reads `_delta_log` history), so disable it on every commit we make.
DISABLE_AUTO_LOG_CLEANUP = deltalake.PostCommitHookProperties(cleanup_expired_logs=False)


def _delta_merge_spill_kwargs() -> dict[str, int]:
"""delta-rs `merge` kwargs that let DataFusion spill to disk instead of OOMing on large merges.

Expand Down Expand Up @@ -151,6 +162,7 @@ def _write_deltalake(
mode=mode,
schema_mode=schema_mode,
commit_properties=commit_properties,
post_commithook_properties=DISABLE_AUTO_LOG_CLEANUP,
)


Expand Down Expand Up @@ -626,6 +638,7 @@ def _do_merge(
predicate=predicate,
streamed_exec=True,
commit_properties=merge_commit_properties,
post_commithook_properties=DISABLE_AUTO_LOG_CLEANUP,
**_delta_merge_spill_kwargs(),
)
.when_matched_update_all()
Expand All @@ -652,6 +665,7 @@ def _do_merge_unpartitioned(data: pa.Table, predicate_ops: list[str]):
predicate=" AND ".join(predicate_ops),
streamed_exec=False,
commit_properties=commit_properties,
post_commithook_properties=DISABLE_AUTO_LOG_CLEANUP,
**_delta_merge_spill_kwargs(),
)
.when_matched_update_all()
Expand Down Expand Up @@ -830,6 +844,7 @@ def _do_scd2_close(first_per_pk: pa.Table, predicate: str) -> dict:
target_alias="target",
predicate=predicate,
streamed_exec=False,
post_commithook_properties=DISABLE_AUTO_LOG_CLEANUP,
**_delta_merge_spill_kwargs(),
)
.when_matched_update(updates={"valid_to": "source.valid_from"})
Expand Down Expand Up @@ -861,6 +876,7 @@ def _do_scd2_close(first_per_pk: pa.Table, predicate: str) -> dict:
mode="append",
schema_mode="merge",
commit_properties=commit_properties,
post_commithook_properties=DISABLE_AUTO_LOG_CLEANUP,
)

delta_table = await self.get_delta_table()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
from products.warehouse_sources.backend.temporal.data_imports.pipelines.core.consts import PARTITION_KEY
from products.warehouse_sources.backend.temporal.data_imports.pipelines.core.delta_table_helper import (
DELTA_MERGE_CONFLICT_RETRIES,
DISABLE_AUTO_LOG_CLEANUP,
DeltaTableHelper,
_delta_merge_spill_kwargs,
_first_per_pk_table,
Expand Down Expand Up @@ -574,6 +575,86 @@ async def test_retries_compact_on_commit_conflict_then_succeeds(self, helper: De
mock_delta.update_incremental.assert_called_once()


class TestPostCommitHookDisablesLogCleanup:
"""Every delta-rs write/merge must pass post_commithook_properties=DISABLE_AUTO_LOG_CLEANUP.

delta-rs's default post-commit hook batch-deletes expired _delta_log files via the same
object_store DeleteObjects call our storage backend can answer with a shape delta-rs's XML
parser rejects ("unknown variant `Code`, expected `Deleted` or `Error`"), failing an
otherwise-successful commit. Regression test for that failure mode reappearing on any of
these call sites.
"""

@pytest.mark.asyncio
async def test_write_disables_log_cleanup(self):
helper = DeltaTableHelper(resource_name="t", job=MagicMock(), logger=_make_logger())
data = pa.table({"id": [1, 2, 3]})
mock_delta = MagicMock()

with (
patch.object(helper, "get_delta_table", AsyncMock(return_value=mock_delta)),
patch.object(helper, "_evolve_delta_schema", AsyncMock(return_value=mock_delta)),
patch("deltalake.write_deltalake") as mock_write,
):
await helper.write_to_deltalake(
data=data, write_type="full_refresh", should_overwrite_table=False, primary_keys=None
)

assert mock_write.call_args.kwargs["post_commithook_properties"] is DISABLE_AUTO_LOG_CLEANUP

@parameterized.expand([("unpartitioned", False), ("partitioned", True)])
@pytest.mark.asyncio
async def test_incremental_merge_disables_log_cleanup(self, _name: str, partitioned: bool):
helper = DeltaTableHelper(resource_name="t", job=MagicMock(), logger=_make_logger())
helper._is_first_sync = False
data_dict: dict[str, Any] = {"id": pa.array([1, 2])}
if partitioned:
data_dict[PARTITION_KEY] = pa.array(["p0", "p0"])
data = pa.table(data_dict)

merge_builder = MagicMock()
merge_builder.when_matched_update_all.return_value = merge_builder
merge_builder.when_not_matched_insert_all.return_value = merge_builder
merge_builder.execute.return_value = {}

mock_delta = MagicMock()
mock_delta.schema.return_value = pa.schema([pa.field("id", pa.int64())])
mock_delta.metadata.return_value = MagicMock(partition_columns=[PARTITION_KEY] if partitioned else [])
mock_delta.merge = MagicMock(return_value=merge_builder)

with (
patch.object(helper, "get_delta_table", AsyncMock(return_value=mock_delta)),
patch.object(helper, "_evolve_delta_schema", AsyncMock(return_value=mock_delta)),
):
await helper.write_to_deltalake(
data=data, write_type="incremental", should_overwrite_table=False, primary_keys=["id"]
)

assert mock_delta.merge.call_args.kwargs["post_commithook_properties"] is DISABLE_AUTO_LOG_CLEANUP

@pytest.mark.asyncio
async def test_scd2_disables_log_cleanup(self):
helper = DeltaTableHelper(resource_name="t", job=MagicMock(), logger=_make_logger())
data = pa.table({"id": [1, 2], "valid_from": [datetime(2026, 1, 1), datetime(2026, 1, 2)]})

close_builder = MagicMock()
close_builder.when_matched_update.return_value = close_builder
close_builder.execute.return_value = {}

mock_delta = MagicMock()
mock_delta.merge = MagicMock(return_value=close_builder)

with (
patch.object(helper, "get_delta_table", AsyncMock(return_value=mock_delta)),
patch.object(helper, "_evolve_delta_schema", AsyncMock(return_value=mock_delta)),
patch("deltalake.write_deltalake") as mock_write,
):
await helper.write_scd2_to_deltalake(data=data, primary_keys=["id"])

assert mock_delta.merge.call_args.kwargs["post_commithook_properties"] is DISABLE_AUTO_LOG_CLEANUP
assert mock_write.call_args.kwargs["post_commithook_properties"] is DISABLE_AUTO_LOG_CLEANUP


def _create_legacy_delta_table(path: str, *, partitioned: bool = False) -> deltalake.DeltaTable:
"""Seed a Delta table that mimics what the old dlt pipeline created:
business columns plus NOT NULL _dlt_id and _dlt_load_id."""
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1741,6 +1741,7 @@ async def test_delta_no_merging_on_first_sync(team, postgres_config, postgres_co
"data": mock.ANY,
"partition_by": mock.ANY,
"commit_properties": mock.ANY,
"post_commithook_properties": mock.ANY,
}

assert second_call_kwargs == {
Expand All @@ -1750,6 +1751,7 @@ async def test_delta_no_merging_on_first_sync(team, postgres_config, postgres_co
"data": mock.ANY,
"partition_by": mock.ANY,
"commit_properties": mock.ANY,
"post_commithook_properties": mock.ANY,
}
else:
mock_v3_post_load.assert_called_once()
Expand All @@ -1766,6 +1768,7 @@ async def test_delta_no_merging_on_first_sync(team, postgres_config, postgres_co
"data": mock.ANY,
"partition_by": mock.ANY,
"commit_properties": mock.ANY,
"post_commithook_properties": mock.ANY,
}

assert second_call_kwargs == {
Expand All @@ -1775,6 +1778,7 @@ async def test_delta_no_merging_on_first_sync(team, postgres_config, postgres_co
"data": mock.ANY,
"partition_by": mock.ANY,
"commit_properties": mock.ANY,
"post_commithook_properties": mock.ANY,
}


Expand Down Expand Up @@ -1842,6 +1846,7 @@ async def test_delta_no_merging_on_first_sync_uncapped_chunk_size(
"data": mock.ANY,
"partition_by": mock.ANY,
"commit_properties": mock.ANY,
"post_commithook_properties": mock.ANY,
}


Expand Down Expand Up @@ -1925,6 +1930,7 @@ async def test_delta_no_merging_on_first_sync_after_reset(team, postgres_config,
"data": mock.ANY,
"partition_by": mock.ANY,
"commit_properties": mock.ANY,
"post_commithook_properties": mock.ANY,
}

assert second_call_kwargs == {
Expand All @@ -1934,6 +1940,7 @@ async def test_delta_no_merging_on_first_sync_after_reset(team, postgres_config,
"data": mock.ANY,
"partition_by": mock.ANY,
"commit_properties": mock.ANY,
"post_commithook_properties": mock.ANY,
}
else:
mock_v3_post_load.assert_called_once()
Expand All @@ -1950,6 +1957,7 @@ async def test_delta_no_merging_on_first_sync_after_reset(team, postgres_config,
"data": mock.ANY,
"partition_by": mock.ANY,
"commit_properties": mock.ANY,
"post_commithook_properties": mock.ANY,
}

assert second_call_kwargs == {
Expand All @@ -1959,6 +1967,7 @@ async def test_delta_no_merging_on_first_sync_after_reset(team, postgres_config,
"data": mock.ANY,
"partition_by": mock.ANY,
"commit_properties": mock.ANY,
"post_commithook_properties": mock.ANY,
}


Expand Down Expand Up @@ -2593,6 +2602,7 @@ async def test_partition_folders_delta_merge_called_with_partition_predicate(
"predicate": f"source.id = target.id AND source.{PARTITION_KEY} = target.{PARTITION_KEY} AND target.{PARTITION_KEY} = '0'",
"streamed_exec": True,
"commit_properties": mock.ANY,
"post_commithook_properties": mock.ANY,
}


Expand Down
Loading