Spark Java数据重复执行问题排查及修复请求
问题分析与解决办法
核心问题拆解
你碰到的每条记录被重复处理4次、transformBlobToString执行22824次的问题,是两个关键问题叠加导致的:
- JDBC分区逻辑错误:用ROWNUM作为分区列,导致每个Spark分区都读取了全量数据
- 多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
相关产品推荐
相关产品推荐

