apache / apache/hudi

[SUPPORT] Queries are very memory intensive due to low read parallelism in HoodieMergeOnReadRDD

Open
#12,434 8 comments 0 reactions 0 assignees View on GitHub
area:performance area:reader
Dominant language
Java
Stars
6.2k
Forks
2.5k
Avg merge
2d 8h
Merged PRs (30d)
111

Description

**Describe the problem you faced**

We have jobs that read from a MOR table using the following pyspark pseudo-code (`event_table_rt` is the MOR table):
```
partitions = ["2023-11-13", "2023-11-14", "2023-11-15", "2023-11-16", "2023-11-17"]
event_df = spark.sql("select * from event_table_rt").filter(
F.col("dt").isin(partitions)
)
user_df = spark.read.format("csv").option("header", "true").load(users_path)
filtered_events_df = df.join(
F.broadcast(user_df),
on=df["user_id"] == user_df["id"],
how="inner",
)
filtered_events_df.write.format("parquet").save("s3://...")
```

We're running into a bottleneck on `HoodieMergeOnReadRDD` (https://github.com/apache/hudi/blob/release-0.14.2/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieMergeOnReadRDD.scala#L37) where the number of tasks in the stage reading `event_df` seems to be non-configurable and (I think) equal to the number of files being read. This is causing massive disk/memory spill and bottlenecking performance.

Is it possible to configure the read parallelism to be higher or is this a fundamental limitation of Hudi with MOR tables? What is the recommendation for how to tune resourcing for readers of MOR tables?

**Environment Description**

* Hudi version : 0.14.1-amzn-1 (EMR 7.2.0)

* Spark version : 3.5.1

* Hive version : 3.1.3

* Hadoop version : 3.3.6

* Storage (HDFS/S3/GCS..) : S3

* Running on Docker? (yes/no) : no

Contributor guide

No contributing guide indexed for this repository

Research direction

Start with hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieMergeOnReadRDD.scala and reproduce the described MOR read using the supplied Spark and Hudi versions. Trace how tasks are selected for the filtered partitions and identify whether a supported parallelism setting exists. Done means a verified configuration or limitation, with the recommended tuning documented or covered by a focused test.

Written by the indexing model from the issue text.

Assessment

Tech stack
python, scala, spark
Domain
data-engineering, distributed-systems, performance
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
24/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.