:wave: question around memory usage for joins + Ra...
# general
g
👋 question around memory usage for joins + Ray, is any memory on the head node used during joins? I was trying to join 10 tables followed by write to parquet and started hitting OOM on the head node even though I am not doing any operation that should collect data to the head node. I have not encountered Ray head node OOM in the past and having trouble figuring out what is being collected there.
c
One cause of head node OOM is that there's a lot of partitions / object refs, which take up memory on the head node for bookkeeping purposes. Do you know roughly how many partitions are going into each join / in total?
💡 1
g
I believe there are a total of 12 joins, each side has 200 partitions
c
Gotcha, one easy thing that may help is
daft.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)
g
one thing to note, the joins are somewhat independent and being concatenated, logic is something like (context: entity resolution between k datasets):
Copy code
def _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
        ]
as number of source datasets grows, this makes the number of joins grow quickly. I am considering just concatenating the source datasets, running a single join and filtering out within data source pairs unless I really want the source matches. or I write out the result
make_candidate_pairs
for each dataset pair separately as opposed to concatenating.
👍 1