From 4091a62e1ffbdabe8fd4ee6b218949f8c5d0235b Mon Sep 17 00:00:00 2001 From: Kevin Yan Date: Wed, 15 Apr 2026 20:34:07 +0000 Subject: [PATCH 1/6] Handle compacting an empty delta --- .../compactor_v2/compaction_session.py | 38 ++++++++++++------- .../compactor_v2/private/compaction_utils.py | 25 +++++++----- 2 files changed, 39 insertions(+), 24 deletions(-) diff --git a/deltacat/compute/compactor_v2/compaction_session.py b/deltacat/compute/compactor_v2/compaction_session.py index c34d15ff1..5cf5163a5 100644 --- a/deltacat/compute/compactor_v2/compaction_session.py +++ b/deltacat/compute/compactor_v2/compaction_session.py @@ -33,6 +33,7 @@ from deltacat.compute.compactor.model.compact_partition_params import ( CompactPartitionParams, ) +from deltacat.utils.common import current_time_ms from deltacat.utils.resources import ( get_current_process_peak_memory_usage_in_bytes, ) @@ -171,27 +172,36 @@ def _execute_compaction( merge_result_list.extend(_run_hash_and_merge_result) # process merge results process_merge_results: tuple[ - Delta, list[MaterializeResult], dict + Optional[Delta], list[MaterializeResult], dict ] = _process_merge_results(params, merge_result_list, compaction_audit) merged_delta, mat_results, hb_id_to_entry_indices_range = process_merge_results - # Record information, logging, and return ExecutionCompactionResult - record_info_msg: str = f" Materialized records: {merged_delta.meta.record_count}" - logger.info(record_info_msg) - compacted_delta: Delta = params.deltacat_storage.commit_delta( - merged_delta, - properties=kwargs.get("properties", {}), - **params.deltacat_storage_kwargs, - ) - logger.info(f"Committed compacted delta: {compacted_delta}") + if merged_delta: + # Record information, logging, and return ExecutionCompactionResult + record_info_msg: str = f" Materialized records: {merged_delta.meta.record_count}" + logger.info(record_info_msg) + compacted_delta: Delta = params.deltacat_storage.commit_delta( + merged_delta, + properties=kwargs.get("properties", {}), + **params.deltacat_storage_kwargs, + ) + new_compacted_delta_locator: DeltaLocator = DeltaLocator.of( + compacted_partition.locator, + compacted_delta.stream_position, + ) + + logger.info(f"Committed compacted delta: {compacted_delta}") + else: + # Avoid committing an empty delta + logger.info("No compacted delta found, skipping committing step...") + new_compacted_delta_locator: DeltaLocator = DeltaLocator.of( + compacted_partition.locator, + current_time_ms(), + ) compaction_end_time: float = time.monotonic() compaction_audit.set_compaction_time_in_seconds( compaction_end_time - compaction_start_time ) - new_compacted_delta_locator: DeltaLocator = DeltaLocator.of( - compacted_partition.locator, - compacted_delta.stream_position, - ) pyarrow_write_result: PyArrowWriteResult = PyArrowWriteResult.union( [m.pyarrow_write_result for m in mat_results] ) diff --git a/deltacat/compute/compactor_v2/private/compaction_utils.py b/deltacat/compute/compactor_v2/private/compaction_utils.py index e9dca6ab3..4e0ed9a47 100644 --- a/deltacat/compute/compactor_v2/private/compaction_utils.py +++ b/deltacat/compute/compactor_v2/private/compaction_utils.py @@ -350,9 +350,10 @@ def _run_hash_and_merge( mutable_compaction_audit.set_records_deduped(total_dd_record_count.item()) mutable_compaction_audit.set_records_deleted(total_deleted_record_count.item()) record_info_msg: str = ( - f"Hash bucket records: {total_hb_record_count}," - f" Deduped records: {total_dd_record_count}, " - f" Deleted records: {total_deleted_record_count}, " + f"Hash bucket records: {total_hb_record_count}, " + f"Input records: {total_input_records_count}, " + f"Deduped records: {total_dd_record_count}, " + f"Deleted records: {total_deleted_record_count}." ) logger.info(record_info_msg) telemetry_this_round = telemetry_time_hb + telemetry_time_merge @@ -560,7 +561,7 @@ def _process_merge_results( params: CompactPartitionParams, merge_results: List[MergeResult], mutable_compaction_audit: CompactionSessionAuditInfo, -) -> tuple[Delta, List[MaterializeResult], dict]: +) -> tuple[Optional[Delta], List[MaterializeResult], dict]: mat_results = [] for merge_result in merge_results: mat_results.extend(merge_result.materialize_results) @@ -602,12 +603,16 @@ def _process_merge_results( **params.s3_client_kwargs, ) deltas: List[Delta] = [m.delta for m in mat_results] - # Note: An appropriate last stream position must be set - # to avoid correctness issue. - merged_delta: Delta = Delta.merge_deltas( - deltas, - stream_position=params.last_stream_position_to_compact, - ) + if deltas: + # Note: An appropriate last stream position must be set + # to avoid correctness issue. + merged_delta: Optional[Delta] = Delta.merge_deltas( + deltas, + stream_position=params.last_stream_position_to_compact, + ) + else: + logger.info("No deltas found to merge!") + merged_delta = None return merged_delta, mat_results, hb_id_to_entry_indices_range From 35c6127ff5f6c535a9e32832712c46eab3f0647d Mon Sep 17 00:00:00 2001 From: Kevin Yan Date: Wed, 15 Apr 2026 20:49:59 +0000 Subject: [PATCH 2/6] lint --- deltacat/compute/compactor_v2/compaction_session.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/deltacat/compute/compactor_v2/compaction_session.py b/deltacat/compute/compactor_v2/compaction_session.py index 5cf5163a5..b3fd5bda7 100644 --- a/deltacat/compute/compactor_v2/compaction_session.py +++ b/deltacat/compute/compactor_v2/compaction_session.py @@ -178,7 +178,9 @@ def _execute_compaction( if merged_delta: # Record information, logging, and return ExecutionCompactionResult - record_info_msg: str = f" Materialized records: {merged_delta.meta.record_count}" + record_info_msg: str = ( + f" Materialized records: {merged_delta.meta.record_count}" + ) logger.info(record_info_msg) compacted_delta: Delta = params.deltacat_storage.commit_delta( merged_delta, From 7792c1a259bbe7442ed900f7ea70b60e2290dd14 Mon Sep 17 00:00:00 2001 From: Kevin Yan Date: Thu, 16 Apr 2026 22:34:26 +0000 Subject: [PATCH 3/6] Commit empty delta as well --- .../compactor_v2/compaction_session.py | 41 ++++++++++++------- 1 file changed, 26 insertions(+), 15 deletions(-) diff --git a/deltacat/compute/compactor_v2/compaction_session.py b/deltacat/compute/compactor_v2/compaction_session.py index b3fd5bda7..cec541aaa 100644 --- a/deltacat/compute/compactor_v2/compaction_session.py +++ b/deltacat/compute/compactor_v2/compaction_session.py @@ -27,7 +27,10 @@ from deltacat.storage import ( Delta, DeltaLocator, + DeltaType, Manifest, + ManifestEntryList, + ManifestMeta, Partition, ) from deltacat.compute.compactor.model.compact_partition_params import ( @@ -182,24 +185,32 @@ def _execute_compaction( f" Materialized records: {merged_delta.meta.record_count}" ) logger.info(record_info_msg) - compacted_delta: Delta = params.deltacat_storage.commit_delta( - merged_delta, - properties=kwargs.get("properties", {}), - **params.deltacat_storage_kwargs, - ) - new_compacted_delta_locator: DeltaLocator = DeltaLocator.of( - compacted_partition.locator, - compacted_delta.stream_position, - ) - - logger.info(f"Committed compacted delta: {compacted_delta}") else: # Avoid committing an empty delta - logger.info("No compacted delta found, skipping committing step...") - new_compacted_delta_locator: DeltaLocator = DeltaLocator.of( - compacted_partition.locator, - current_time_ms(), + logger.info("No compacted delta found, committing an empty delta...") + merged_delta = Delta.of( + DeltaLocator.of( + compacted_partition.locator, + None, + ), + DeltaType.UPSERT, + ManifestMeta(), + {}, + Manifest.of(ManifestEntryList()), ) + + compacted_delta: Delta = params.deltacat_storage.commit_delta( + merged_delta, + properties=kwargs.get("properties", {}), + **params.deltacat_storage_kwargs, + ) + new_compacted_delta_locator: DeltaLocator = DeltaLocator.of( + compacted_partition.locator, + compacted_delta.stream_position, + ) + + logger.info(f"Committed compacted delta: {compacted_delta}") + compaction_end_time: float = time.monotonic() compaction_audit.set_compaction_time_in_seconds( compaction_end_time - compaction_start_time From 7373c66cb17a0b9c46c52b9598c6af785fa64cdd Mon Sep 17 00:00:00 2001 From: Kevin Yan Date: Thu, 16 Apr 2026 22:44:02 +0000 Subject: [PATCH 4/6] Create empty delta within compaction_utils --- .../compactor_v2/compaction_session.py | 42 ++++++------------- .../compactor_v2/private/compaction_utils.py | 21 ++++++++-- 2 files changed, 29 insertions(+), 34 deletions(-) diff --git a/deltacat/compute/compactor_v2/compaction_session.py b/deltacat/compute/compactor_v2/compaction_session.py index cec541aaa..e9a80d6d0 100644 --- a/deltacat/compute/compactor_v2/compaction_session.py +++ b/deltacat/compute/compactor_v2/compaction_session.py @@ -27,16 +27,12 @@ from deltacat.storage import ( Delta, DeltaLocator, - DeltaType, Manifest, - ManifestEntryList, - ManifestMeta, Partition, ) from deltacat.compute.compactor.model.compact_partition_params import ( CompactPartitionParams, ) -from deltacat.utils.common import current_time_ms from deltacat.utils.resources import ( get_current_process_peak_memory_usage_in_bytes, ) @@ -160,6 +156,7 @@ def _execute_compaction( logger.info(f"Length of grouped uniform deltas is: {len(uniform_deltas_grouped)}") merge_result_list: List[MergeResult] = [] compacted_partition = _stage_new_partition(params) + logger.info(f"Got staged compacted partition: {compacted_partition}") for uniform_deltas in uniform_deltas_grouped: # run hash and merge _run_hash_and_merge_result: List[MergeResult] = _run_hash_and_merge( @@ -175,46 +172,31 @@ def _execute_compaction( merge_result_list.extend(_run_hash_and_merge_result) # process merge results process_merge_results: tuple[ - Optional[Delta], list[MaterializeResult], dict - ] = _process_merge_results(params, merge_result_list, compaction_audit) + Delta, list[MaterializeResult], dict + ] = _process_merge_results( + params, merge_result_list, compaction_audit, compacted_partition + ) merged_delta, mat_results, hb_id_to_entry_indices_range = process_merge_results - if merged_delta: - # Record information, logging, and return ExecutionCompactionResult - record_info_msg: str = ( - f" Materialized records: {merged_delta.meta.record_count}" - ) - logger.info(record_info_msg) - else: - # Avoid committing an empty delta - logger.info("No compacted delta found, committing an empty delta...") - merged_delta = Delta.of( - DeltaLocator.of( - compacted_partition.locator, - None, - ), - DeltaType.UPSERT, - ManifestMeta(), - {}, - Manifest.of(ManifestEntryList()), - ) + # Record information, logging, and return ExecutionCompactionResult + record_info_msg: str = f" Materialized records: {merged_delta.meta.record_count}" + logger.info(record_info_msg) compacted_delta: Delta = params.deltacat_storage.commit_delta( merged_delta, properties=kwargs.get("properties", {}), **params.deltacat_storage_kwargs, ) - new_compacted_delta_locator: DeltaLocator = DeltaLocator.of( - compacted_partition.locator, - compacted_delta.stream_position, - ) logger.info(f"Committed compacted delta: {compacted_delta}") - compaction_end_time: float = time.monotonic() compaction_audit.set_compaction_time_in_seconds( compaction_end_time - compaction_start_time ) + new_compacted_delta_locator: DeltaLocator = DeltaLocator.of( + compacted_partition.locator, + compacted_delta.stream_position, + ) pyarrow_write_result: PyArrowWriteResult = PyArrowWriteResult.union( [m.pyarrow_write_result for m in mat_results] ) diff --git a/deltacat/compute/compactor_v2/private/compaction_utils.py b/deltacat/compute/compactor_v2/private/compaction_utils.py index 4e0ed9a47..009389799 100644 --- a/deltacat/compute/compactor_v2/private/compaction_utils.py +++ b/deltacat/compute/compactor_v2/private/compaction_utils.py @@ -49,6 +49,8 @@ DeltaLocator, Partition, Manifest, + ManifestEntryList, + ManifestMeta, Stream, StreamLocator, ) @@ -561,7 +563,8 @@ def _process_merge_results( params: CompactPartitionParams, merge_results: List[MergeResult], mutable_compaction_audit: CompactionSessionAuditInfo, -) -> tuple[Optional[Delta], List[MaterializeResult], dict]: + compacted_partition: Partition, +) -> tuple[Delta, List[MaterializeResult], dict]: mat_results = [] for merge_result in merge_results: mat_results.extend(merge_result.materialize_results) @@ -606,13 +609,23 @@ def _process_merge_results( if deltas: # Note: An appropriate last stream position must be set # to avoid correctness issue. - merged_delta: Optional[Delta] = Delta.merge_deltas( + merged_delta: Delta = Delta.merge_deltas( deltas, stream_position=params.last_stream_position_to_compact, ) else: - logger.info("No deltas found to merge!") - merged_delta = None + logger.info("No deltas found to merge, returning an empty delta...") + empty_manifest = Manifest.of(ManifestEntryList()) + merged_delta: Delta = Delta.of( + DeltaLocator.of( + compacted_partition.locator, + None, + ), + DeltaType.UPSERT, + empty_manifest.meta, + {}, + empty_manifest, + ) return merged_delta, mat_results, hb_id_to_entry_indices_range From 19e5a48ec8d870c5be6831bed8372fd050f0f82e Mon Sep 17 00:00:00 2001 From: Kevin Yan Date: Thu, 16 Apr 2026 23:34:24 +0000 Subject: [PATCH 5/6] lint --- deltacat/compute/compactor_v2/private/compaction_utils.py | 1 - 1 file changed, 1 deletion(-) diff --git a/deltacat/compute/compactor_v2/private/compaction_utils.py b/deltacat/compute/compactor_v2/private/compaction_utils.py index 009389799..61b2364b1 100644 --- a/deltacat/compute/compactor_v2/private/compaction_utils.py +++ b/deltacat/compute/compactor_v2/private/compaction_utils.py @@ -50,7 +50,6 @@ Partition, Manifest, ManifestEntryList, - ManifestMeta, Stream, StreamLocator, ) From 2a8fa5be78e83a31ffc4be00b1846ad888fe47a8 Mon Sep 17 00:00:00 2001 From: Kevin Yan Date: Fri, 17 Apr 2026 18:06:22 +0000 Subject: [PATCH 6/6] Use current_time_ms() for empty delta stream position --- deltacat/compute/compactor_v2/private/compaction_utils.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/deltacat/compute/compactor_v2/private/compaction_utils.py b/deltacat/compute/compactor_v2/private/compaction_utils.py index 61b2364b1..0638e2c0c 100644 --- a/deltacat/compute/compactor_v2/private/compaction_utils.py +++ b/deltacat/compute/compactor_v2/private/compaction_utils.py @@ -56,6 +56,7 @@ from deltacat.compute.compactor.model.compact_partition_params import ( CompactPartitionParams, ) +from deltacat.utils.common import current_time_ms from deltacat.utils.ray_utils.concurrency import ( invoke_parallel, task_resource_options_provider, @@ -618,13 +619,14 @@ def _process_merge_results( merged_delta: Delta = Delta.of( DeltaLocator.of( compacted_partition.locator, - None, + current_time_ms(), ), DeltaType.UPSERT, empty_manifest.meta, {}, empty_manifest, ) + logger.info(f"Created empty delta: {merged_delta}") return merged_delta, mat_results, hb_id_to_entry_indices_range