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

Spark Java数据重复执行问题排查及修复请求

问题分析与解决办法

核心问题拆解

你碰到的每条记录被重复处理4次、transformBlobToString执行22824次的问题,是两个关键问题叠加导致的:

  1. JDBC分区逻辑错误:用ROWNUM作为分区列,导致每个Spark分区都读取了全量数据
  2. 多Action触发重复计算:Spark懒执行机制下,多次Action操作会重复执行整个数据处理链路

问题1:JDBC分区用ROWNUM导致全量读取

Oracle的ROWNUM是查询时动态生成的伪列,Spark的JDBC分区会为每个分区生成独立的查询语句。你的SQL写法会让每个分区都先执行select * from myTableName where id=?全量查询,再过滤ROWNUM范围,相当于4个分区各读了一遍全量数据,直接导致每条记录被处理4次。

修正方案

改用表中真实存在的有序列(比如主键、自增ID)作为分区列,示例代码:

public Dataset<Row> getRecords(Long id) throws SQLException {
    String table = "myTableName";
    // 获取该id对应的主键范围(需要你实现getMinPkId和getMaxPkId方法)
    String minPk = String.valueOf(getMinPkId(id));
    String maxPk = String.valueOf(getMaxPkId(id));
    String totalRowCount = getTotalRowCount(id, otherFIlters);

    Properties properties = new Properties();
    properties.setProperty("partitionColumn", "pk_id"); // 替换成你的真实有序列名
    properties.setProperty("lowerBound", minPk);
    properties.setProperty("upperBound", maxPk);
    properties.setProperty("numPartitions", "4");
    properties.setProperty("fetchsize", "100");
    properties.setProperty("Driver", driver);
    properties.setProperty("user", user);
    properties.setProperty("password", password);

    // 直接带where条件查询原表,避免嵌套ROWNUM的问题
    Dataset<Row> records = SparkSession.getActiveSession().get()
        .read()
        .jdbc(jdbcUrl, table, "id = " + id, properties);
    return records;
}

如果必须用ROWNUM,需要先将数据写入临时表,把ROWNUM转为固定列后再分区读取,但这种方式效率较低,优先推荐用真实有序列。

问题2:多Action触发重复计算

你的calculatePartitions(csvDataSet)方法大概率调用了count()这类Action操作,加上后续的write().save()又是一次Action,两次Action会触发整个数据处理链路(从JDBC读取到转换)执行两次。结合问题1的每个分区全量读取,就会让重复处理的情况更严重。

修正方案

在生成csvDataSet后立即缓存,避免重复计算:

// 生成csvDataSet后缓存数据,选择合适的存储级别
csvDataSet = filterFormatDataFrame(xmlDataSet);
csvDataSet.persist(StorageLevel.MEMORY_AND_DISK_SER()); // 内存不足时写入磁盘,序列化节省空间

int requiredPartitions = calculatePartitions(csvDataSet);
if (requiredPartitions > dbSvc.getNumberOfDBPartitions()) {
    csvDataSet.repartition(requiredPartitions)
        .write().format("csv").option("header", "true")
        .save(filepath);
} else {
    csvDataSet.write().format("csv").option("header", "true")
        .save(filepath);
}
// 作业完成后手动释放缓存,避免占用资源
csvDataSet.unpersist();

缓存后,无论触发多少次Action,都只会执行一次全量计算,后续Action直接读取缓存中的数据。

额外检查点

  • 确认processRecords中的transformBlobToString是针对单条记录的处理逻辑,没有在分区内重复执行
  • 验证getTotalRowCount的查询结果准确,避免分区范围计算错误导致数据分布不均

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 02:50:52