Garrett Weaver
01/28/2025, 11:25 PMColin Ho
01/29/2025, 12:00 AMGarrett Weaver
01/29/2025, 12:01 AMColin Ho
01/29/2025, 12:12 AMdaft.context.set_execution_config(shuffle_algorithm="pre_shuffle_merge") , which will try to merge small partitions together during a shuffle.
(Up to 1GB by default, tunable by setting pre_shuffle_merge_threshold in the config)Garrett Weaver
01/29/2025, 12:59 AMdef _find_between_source_edges(
self, source_datasets: List[SourceDataset[DataFrameTypeT]]
) -> List[DataFrameTypeT]:
"""Finds between-source edges / candidate pairs."""
return [
blocking_rule.make_candidate_pairs(
dataset_left.data_renamed, dataset_right.data_renamed
)
for dataset_left, dataset_right in itertools.combinations(
source_datasets, 2
)
for blocking_rule in self._blocking_rules
]Garrett Weaver
01/29/2025, 1:00 AMmake_candidate_pairs for each dataset pair separately as opposed to concatenating.