You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark作业为何执行两次数据库读取?如何避免重复读取?

Why does my Spark workflow read the same database twice, and how to fix it?

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

  1. Spark's lazy execution behavior
    Spark doesn't actually process any data until it hits an action (like save(), count(), or operations that resolve lineage). If your modifiedDf and rollUpDf both 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 generating rollUpDf involves an implicit action, or the final union + save triggers two separate passes over the lineage, you’ll end up with duplicate reads.

  2. Accidental duplicate read logic
    It’s easy to accidentally write code that reads the source twice instead of reusing the same rawDf instance. 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 reading rawDf, call cache() or persist() 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 datasets
    

    This tells Spark to keep the raw data around after the first read, so both modifiedDf and rollUpDf can reuse it.

  • Reuse the same rawDf instance
    Double-check your code to ensure both modifiedDf and rollUpDf are derived from the same rawDf variable, 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-reading
    
  • Verify execution plans with explain()
    Use modifiedDf.explain() and rollUpDf.explain() to inspect the query plan. If you see Scan JDBCRelation (or equivalent for your data source) in both plans, that confirms the duplicate read. After caching rawDf, the plan should show InMemoryTableScan for subsequent transformations instead.

  • Force early materialization (for large datasets)
    If caching alone isn’t feasible (e.g., for extremely large datasets), trigger an action on rawDf right 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 data
    

    This ensures all downstream transformations use the already-loaded data.

Content of the question comes from Stack Exchange, question author: LN.EXE

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.20 12:30:36