Spark Scala优化:为DataFrame添加去重列表列的高效实现
Optimized Solution for Adding Site-Specific Offer Lists
Here's the corrected, efficient approach to achieve your requirement:
// Step 1: Precompute distinct offer lists per site (one-time distributed aggregation) val siteOffers = finvResultOutput1 .groupBy("Site_siteId") .agg(collect_set("OfferSpecification_displayLabel").alias("offersAtSite")) // Step 2: Join back to original DataFrame and apply conditional logic val finalResult = finvResultOutput1 .join(siteOffers, Seq("Site_siteId"), "left") .withColumn("offersAtSite", when(col("activeIOsAtSite") === "Y", col("offersAtSite")).otherwise(null))
Why Your Original Code Fails
- Catastrophic Performance: Using
collect()inside awhenclause pulls the entire dataset onto the driver node, which is drastically slow for large data and often causes out-of-memory crashes. This operation runs once globally, not per row, so it can't deliver site-specific values as intended. - Incorrect Filtering: The subquery's
where(col("Site_siteId") === col("siteid"))is a no-op—you're comparing the column to itself (after aliasing), so it filters nothing. This meanscollect()returns all offer labels from the entire dataset, not just those matching the current row's site ID.
Key Improvements
- Distributed Processing:
groupBy+collect_setruns in parallel across the cluster, avoiding driver bottlenecks.collect_setautomatically handles deduplication, which aligns with your "distinct list" requirement. - Efficient Lookup: Joining the precomputed site-offer mapping to the original DataFrame ensures each row gets the correct list without redundant calculations.
- Precise Conditional Logic: The
whenclause correctly applies the list only whenactiveIOsAtSiteis "Y", setting it tonullotherwise (replacenullwitharray()if you prefer an empty list instead).
内容的提问来源于stack exchange,提问作者mr.Penguin
相关产品推荐
相关产品推荐

