diff --git a/deltacat/compute/compactor_v2/compaction_session.py b/deltacat/compute/compactor_v2/compaction_session.py index c34d15ff1..e9a80d6d0 100644 --- a/deltacat/compute/compactor_v2/compaction_session.py +++ b/deltacat/compute/compactor_v2/compaction_session.py @@ -156,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( @@ -172,11 +173,15 @@ def _execute_compaction( # process merge results process_merge_results: tuple[ Delta, list[MaterializeResult], dict - ] = _process_merge_results(params, merge_result_list, compaction_audit) + ] = _process_merge_results( + params, merge_result_list, compaction_audit, compacted_partition + ) 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", {}), diff --git a/deltacat/compute/compactor_v2/private/compaction_utils.py b/deltacat/compute/compactor_v2/private/compaction_utils.py index e9dca6ab3..0638e2c0c 100644 --- a/deltacat/compute/compactor_v2/private/compaction_utils.py +++ b/deltacat/compute/compactor_v2/private/compaction_utils.py @@ -49,12 +49,14 @@ DeltaLocator, Partition, Manifest, + ManifestEntryList, Stream, StreamLocator, ) 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, @@ -350,9 +352,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,6 +563,7 @@ def _process_merge_results( params: CompactPartitionParams, merge_results: List[MergeResult], mutable_compaction_audit: CompactionSessionAuditInfo, + compacted_partition: Partition, ) -> tuple[Delta, List[MaterializeResult], dict]: mat_results = [] for merge_result in merge_results: @@ -602,12 +606,27 @@ 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: Delta = Delta.merge_deltas( + deltas, + stream_position=params.last_stream_position_to_compact, + ) + else: + 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, + 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