diff --git a/products/warehouse_sources/backend/temporal/data_imports/pipelines/core/delta_table_helper.py b/products/warehouse_sources/backend/temporal/data_imports/pipelines/core/delta_table_helper.py index ea946a7d16e4..26bcfc7e373b 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/pipelines/core/delta_table_helper.py +++ b/products/warehouse_sources/backend/temporal/data_imports/pipelines/core/delta_table_helper.py @@ -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. @@ -151,6 +162,7 @@ def _write_deltalake( mode=mode, schema_mode=schema_mode, commit_properties=commit_properties, + post_commithook_properties=DISABLE_AUTO_LOG_CLEANUP, ) @@ -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() @@ -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() @@ -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"}) @@ -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() diff --git a/products/warehouse_sources/backend/temporal/data_imports/pipelines/core/test/test_delta_table_helper.py b/products/warehouse_sources/backend/temporal/data_imports/pipelines/core/test/test_delta_table_helper.py index 644c6f62feb9..953f9ef00e9b 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/pipelines/core/test/test_delta_table_helper.py +++ b/products/warehouse_sources/backend/temporal/data_imports/pipelines/core/test/test_delta_table_helper.py @@ -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, @@ -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.""" diff --git a/products/warehouse_sources/backend/temporal/data_imports/tests/e2e/test_end_to_end.py b/products/warehouse_sources/backend/temporal/data_imports/tests/e2e/test_end_to_end.py index f365535c48f8..3b84c0bdf768 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/tests/e2e/test_end_to_end.py +++ b/products/warehouse_sources/backend/temporal/data_imports/tests/e2e/test_end_to_end.py @@ -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 == { @@ -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() @@ -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 == { @@ -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, } @@ -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, } @@ -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 == { @@ -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() @@ -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 == { @@ -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, } @@ -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, }