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

Spark Save操作重复执行Map阶段问题咨询

问题描述

业务场景如下:

  1. 从数据库读取4分区的数据:
properties.setProperty("partitionColumn", "num_rows");
properties.setProperty("lowerBound", "0");
properties.setProperty("upperBound", getTotalRowCount(ID));
properties.setProperty("numPartitions", "4");
properties.setProperty("Driver", driver);
properties.setProperty("user", user);
properties.setProperty("password", password);
Dataset<Row> records = SparkSession.getActiveSession().get().read().jdbc(jdbcUrl, table, properties);
  1. 对数据集进行转换处理:
Dataset<String> stringRecordSet = dbRecordsSet.map((MapFunction<Row, String> )xmlRow -> {
        return TransformationService.extractXMLBlobToString(xmlRow);
    }, Encoders.STRING());   
Dataset<Row> jsonDataSet = sparkSession.read().json(stringRecordSet);
  1. 将结果保存为CSV文件:
jsonDataSet.write().format("csv").save(filepath);

遇到的问题:
数据集共40行,步骤2会按预期并行处理各分区,但执行步骤3保存时,会重新执行步骤2的处理逻辑,map中的TransformationService.extractXMLBlobToString方法会被再次调用处理全部40行数据。生产环境中数据集达百万级,会导致数据被不必要地重复处理两次。而如果在步骤2和3之间添加缓存,处理时间反而从5分钟增至10分钟,完全没有帮助。请问为何Save操作会触发Map阶段的重复执行?


问题原因与解决方案

一、重复执行的核心原因

Spark的惰性求值机制是根本原因:

  • Spark中所有转换操作(比如map、filter)都是惰性的,不会立即执行,只会记录执行逻辑形成DAG(有向无环图)。
  • 只有遇到行动操作(比如write.save、count、collect)时,才会触发整个DAG的计算。
  • 你的场景中,sparkSession.read().json(stringRecordSet)是隐式行动操作(读取Dataset生成新Dataset时,需要触发上游计算获取数据),这会第一次触发步骤2的map执行;之后的write.save是第二个行动操作,会再次触发整个DAG重新计算,所以map逻辑被执行了两次。

二、缓存失效的原因

添加缓存后耗时增加,大概率是以下问题:

  1. 存储介质选择不当:默认缓存用内存,如果数据集无法完全放入内存,会溢写到磁盘,磁盘IO会大幅增加耗时。
  2. 缓存位置错误:若只缓存stringRecordSet,但jsonDataSet基于它生成时未复用缓存,或缓存的数据集后续被重新转换,缓存未被有效利用。
  3. 缓存未触发持久化:cache()方法是惰性的,若缓存后未先触发轻量行动操作(比如count()),后续write会先计算再缓存,等于做了两次计算。

三、正确解决方法

  1. 调整缓存策略:
    • 对最终的jsonDataSet做持久化,指定合适的存储级别,比如内存不足时用MEMORY_AND_DISK_SER(序列化后存内存+磁盘),减少内存占用和IO开销:
      jsonDataSet.persist(StorageLevel.MEMORY_AND_DISK_SER());
      
    • 缓存后先触发一次轻量行动操作(比如count()),提前把数据存入缓存,避免后续write重复计算:
      jsonDataSet.count();
      
  2. 避免隐式行动操作:
    • 先缓存map后的stringRecordSet,再生成jsonDataSet,确保read.json时复用缓存数据,不重新计算map逻辑。
  3. 检查DAG依赖:
    • 通过Spark UI查看DAG图,确认是否存在重复计算分支,确保缓存的数据集是后续操作的直接依赖,避免DAG生成重复计算步骤。

内容的提问来源于stack exchange,提问作者Nila

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 14:05:14