Skip to content
Merged
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
7 changes: 6 additions & 1 deletion deltacat/compute/compactor_v2/compaction_session.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand All @@ -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", {}),
Expand Down
37 changes: 28 additions & 9 deletions deltacat/compute/compactor_v2/private/compaction_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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

Expand Down
Loading