Spark作业为何执行两次数据库读取?如何避免重复读取?
Question Description
My data workflow is as follows:
rawDf -> modifiedDf -> rollUpDf -> union(modifiedDf, rollUpDf) -> save. Currently, the performance isn't meeting expectations, and after troubleshooting, I found there are two database read operations, occurring at stage 60 (when generating rollUpDf) and stage 61. Since both modifiedDf and rollUpDf come from the same data source, I don't understand why there are two reads. I'm asking about the reason for this phenomenon and how to avoid it.
Answer
Let me break down why this is happening and how to fix it—this is a super common gotcha with Spark's lazy evaluation model!
Why the double read occurs
Spark's lazy execution behavior
Spark doesn't actually process any data until it hits an action (likesave(),count(), or operations that resolve lineage). If yourmodifiedDfandrollUpDfboth trace their lineage back to the original database read, and you haven't persisted the raw data, Spark will re-run the entire lineage (including the database read) every time an action is triggered. For example, if generatingrollUpDfinvolves an implicit action, or the finalunion + savetriggers two separate passes over the lineage, you’ll end up with duplicate reads.Accidental duplicate read logic
It’s easy to accidentally write code that reads the source twice instead of reusing the samerawDfinstance. For example:// Bad: Reads the database twice val modifiedDf = spark.read.jdbc(...) // First read .transform(...) val rollUpDf = spark.read.jdbc(...) // Second read—unintentional duplicate! .groupBy(...)Even with identical read parameters, Spark treats these as separate operations and hits the database twice.
How to avoid duplicate reads
Here are the most reliable fixes:
Cache or persist your rawDf
Right after readingrawDf, callcache()orpersist()to store the raw data in memory (and disk if needed) so Spark doesn’t re-read the database for subsequent transformations:val rawDf = spark.read.jdbc(...) rawDf.cache() // Use persist(StorageLevel.MEMORY_AND_DISK) for larger datasetsThis tells Spark to keep the raw data around after the first read, so both
modifiedDfandrollUpDfcan reuse it.Reuse the same rawDf instance
Double-check your code to ensure bothmodifiedDfandrollUpDfare derived from the samerawDfvariable, not separate read calls:// Good: Reuses the already-read rawDf val rawDf = spark.read.jdbc(...) val modifiedDf = rawDf.transform(...) val rollUpDf = rawDf.groupBy(...) // Uses the existing rawDf instead of re-readingVerify execution plans with explain()
UsemodifiedDf.explain()androllUpDf.explain()to inspect the query plan. If you seeScan JDBCRelation(or equivalent for your data source) in both plans, that confirms the duplicate read. After cachingrawDf, the plan should showInMemoryTableScanfor subsequent transformations instead.Force early materialization (for large datasets)
If caching alone isn’t feasible (e.g., for extremely large datasets), trigger an action onrawDfright after reading to force the database read to happen once upfront:val rawDf = spark.read.jdbc(...) rawDf.count() // Triggers the initial read and materializes the dataThis ensures all downstream transformations use the already-loaded data.
Content of the question comes from Stack Exchange, question author: LN.EXE

