From b6f396f3ca21049aada5f6ba0680de17cdc22075 Mon Sep 17 00:00:00 2001 From: Tom Owers Date: Thu, 30 Jul 2026 12:15:23 +0100 Subject: [PATCH 1/2] fix(data-imports): disable automatic delta-rs log cleanup on writes/merges Every delta-rs write and merge commit runs a post-commit hook that, per table defaults, batch-deletes expired `_delta_log` JSON files once a commit's history crosses the retention window. That cleanup uses the same object_store bulk `DeleteObjects` call our storage backend can answer with a response shape delta-rs's XML parser doesn't recognize, failing the whole commit even though the write itself already succeeded. We don't rely on this automatic cleanup, so disable it explicitly on every write and merge in `delta_table_helper.py`. Co-Authored-By: PostHog Code Generated-By: PostHog Code Task-Id: 710c9b62-30ba-4ac7-b606-27fce571ea6b --- .../pipelines/core/delta_table_helper.py | 16 ++++ .../core/test/test_delta_table_helper.py | 81 +++++++++++++++++++ 2 files changed, 97 insertions(+) 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.""" From f97757c4e83883b6f5e3d7f7f866dfe973cbdf45 Mon Sep 17 00:00:00 2001 From: Tom Owers Date: Thu, 30 Jul 2026 12:15:25 +0100 Subject: [PATCH 2/2] fix(data-imports): update e2e delta write/merge kwarg assertions for new post_commithook_properties kwarg CI caught this: the previous commit added post_commithook_properties to every delta-rs write/merge call, but several e2e tests assert exact-equality on the kwargs dict passed to write_deltalake/merge, so they failed on the new key. Add it to each expected dict. Generated-By: PostHog Code Task-Id: 710c9b62-30ba-4ac7-b606-27fce571ea6b --- .../temporal/data_imports/tests/e2e/test_end_to_end.py | 10 ++++++++++ 1 file changed, 10 insertions(+) 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, }