【问题标题】:Spark Error - Max iterations (100) reached for batch ResolutionSpark 错误 - 达到批量解决的最大迭代次数 (100)
【发布时间】:2020-04-09 15:24:59
【问题描述】:

我正在研究 Spark SQL,我需要找出两个大型 CSV 之间的差异。

Diff 应该给出:-

  • 插入的行或新记录 // 仅比较 Id 的

  • 更改的行(不包括插入的行)- 比较所有列值

  • 已删除的行 // 仅比较 Id 的

Spark 2.4.4 + Java

我正在使用 Databricks 读取/写入 CSV

Dataset<Row> insertedDf = newDf_temp.join(oldDf_temp,oldDf_temp.col(key)
                .equalTo(newDf_temp.col(key)),"left_anti");
Long insertedCount = insertedDf.count();
logger.info("Inserted File Count == "+insertedCount);


Dataset<Row> deletedDf = oldDf_temp.join(newDf_temp,oldDf_temp.col(key)
                .equalTo(newDf_temp.col(key)),"left_anti")
                .select(oldDf_temp.col(key));
Long deletedCount = deletedDf.count();
logger.info("deleted File Count == "+deletedCount);


Dataset<Row> changedDf = newDf_temp.exceptAll(oldDf_temp); // This gives rows (New +changed Records)

Dataset<Row> changedDfTemp = changedDf.join(insertedDf, changedDf.col(key)
                .equalTo(insertedDf.col(key)),"left_anti"); // This gives only changed record

Long changedCount = changedDfTemp.count();
logger.info("Changed File Count == "+changedCount);

这适用于最多 50 列左右的 CSV。

The Above code fails for one row in CSV with 300+columns, so I am sure this is not file Size problem.

但如果我有一个包含 300 多列的 CSV,那么它会因异常而失败

达到批量解决的最大迭代次数 (100) - Spark 错误

If I set the below property in Spark, It Works!!!

sparkConf.set("spark.sql.optimizer.maxIterations", "500");

但我的问题是为什么我必须设置这个?

我做错了什么吗? 或者对于具有大列的 CSV,这种行为是预期的。

我能否以任何方式对其进行优化以处理大列 CSV。

【问题讨论】:

    标签: apache-spark apache-spark-sql data-science


    【解决方案1】:

    您遇到的问题与 spark 如何接受您告诉它的指令并将其转换为它要执行的实际操作有关。它首先需要通过运行 Analyzer 来理解您的指令,然后它会尝试通过运行其优化器来改进它们。该设置似乎适用于两者。

    具体来说,您的代码在分析器中的某个步骤中被炸毁。分析器负责确定您何时引用事物,您实际指的是什么事物。例如,将函数名称映射到实现或跨重命名和不同转换映射列名称。它在多次传递中执行此操作,每次传递都解决额外的问题,然后再次检查它是否可以解决移动。

    我认为您的情况发生的情况是每次通过可能会解决一列,但 100 次通过不足以解决所有列。通过增加它,您可以为其提供足够的通行证,以便能够完全完成您的计划。对于潜在的性能问题,这绝对是一个危险信号,但如果您的代码正在运行,那么您可能只需增加该值而不必担心它。

    如果它不起作用,那么您可能需要尝试做一些事情来减少计划中使用的列数。也许将所有列组合成一个编码字符串列作为键。在进行联接之前检查点数据可能会使您受益,因此您可以缩短计划。

    编辑:

    另外,我会重构您上面的代码,以便您只需一个连接即可完成所有操作。这应该会快很多,并且可能会解决您的其他问题。

    每次加入都会导致随机播放(数据在计算节点之间发送),这会增加您的工作时间。无需独立计算添加、删除和更改,您可以一次完成所有操作。类似于下面的代码。它在 scala 伪代码中,因为我比 Java API 更熟悉它。

    import org.apache.spark.sql.functions._
    
    var oldDf = ..
    var newDf = ..
    val changeCols = newDf.columns.filter(_ != "id").map(col)
    
    // Make the columns you want to compare into a single struct column for easier comparison
    newDf = newDF.select($"id", struct(changeCols:_*) as "compare_new")
    oldDf = oldDF.select($"id", struct(changeCols:_*) as "compare_old")
    
    // Outer join on ID
    val combined = oldDF.join(newDf, Seq("id"), "outer")
    
    // Figure out status of each based upon presence of old/new
    //  IF old side is missing, must be an ADD
    //  IF new side is missing, must be a DELETE
    //  IF both sides present but different, it's a CHANGE
    //  ELSE it's NOCHANGE
    val status = when($"compare_new".isNull, lit("add")).
                 when($"compare_old".isNull, lit("delete")).
                 when($"$compare_new" != $"compare_old", lit("change")).
                 otherwise(lit("nochange"))
    
    val labeled = combined.select($"id", status)
    

    此时,我们将每个 ID 都标记为 ADD/DELETE/CHANGE/NOCHANGE,因此我们可以只使用 groupBy/count。这个 agg 几乎可以完全在 map 端完成,所以它会比 join 快很多。

    labeled.groupBy("status").count.show
    

    【讨论】:

    • 感谢您的回复,因为我是新来的火花,这是正常的吗?我知道 Spark 是否可以处理大量数据。我将搜索一些示例来组合列或检查点,看看是否有帮助。
    • 并不常见。问题不在于数据的大小(以计数/字节计),而在于数据的宽度。即使这样通常也很好,因为您通常不需要触摸数据集中的每一列,但是您想要将所有 300 列领导用于分析器的事实需要比平时做更多的工作。
    • 我添加了一个可能有帮助的重构建议。这是构建 Spark 代码的更好方法,它可能会解决您的分析器问题,但不确定。
    猜你喜欢
    • 2013-06-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-03-06
    • 2016-02-25
    • 2015-05-07
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多