Spark process data more than available memory

Viewed 42

I'm working with Apache Spark version 3.1.2 deployed on a cluster of 4 nodes, each having 24GB of memory and 8 cores i.e. ~96GB of distributed memory. I want to read-in and process about ~120GB of compressed (gzip) json data.

Following is a generic code flow of my processing

        data = spark.read.option('multiline', True).json(data_path, schema=schema)
        result = data.filter(data['col_1']['col_1_1'].isNotNull() | data['col2'].isNotNull()) \
                     .rdd \
                     .map(parse_json_and_select_columns_of_interest) \
                     .toDF(schema_of_interest) \
                     .filter(data['col_x'].isin(broadcast_filter_list)) \
                     .rdd \
                     .map(lambda x: (x['col_key'], x.asDict())) \
                     .groupByKey() \
                     .mapValues(compute_and_add_extra_columns) \
                     .flatMap(...) \
                     .reduceByKey(lambda a,b:a+b) \  <--- OOM
                     .sortByKey()
                     .map(append_columns_based_on_key)
                     .saveAsTextFile(...)

I have tried with following executor settings

# Tiny executors
--num-executors 32
--executor-cores 1
--executor-memory 2g
# Fat executors
--num-executors 4
--executor-cores 8
--executor-memory 20g

However, for all of these settings, I keep getting out of memory especially on .reduceByKey(lambda a,b:a+b). My question is, (1) Regardless of performance, can I change my code flow to avoid getting OOM? or (2) Should I add more memory to my cluster? (Avoiding this since it may not be a sustainable solution in long run)

Thanks

1 Answers

I would actually guess it's the sortByKey causing the OOM and would suggest increasing the number of partitions you are using by passing an argument sortByKey(numPartitions = X).

Also, I can suggest trying to use DataFrame API where possible.

Related