【问题标题】:Spark accumulator not counting correctly?火花蓄电池计数不正确?
【发布时间】:2017-05-26 10:53:06
【问题描述】:

使用 Spark 2.1,我有一个函数,它接受 DataFrame 并检查是否所有记录都在给定的数据库中(在本例中为 Aerospike)。

看起来很像这样:

def check(df: DataFrame): Long = {
    val finalResult = df.sparkSession.sparkContext.longAccumulator("finalResult")
    df.rdd.foreachPartition(iter => {
        val success = //if record is on the database: 1 else: 0 
        //if success = 0, send Slack message with missing record
        finalResult.add(success)
       }
      df.count - finalResult.value
    }

因此,Slack 消息的数量应该与函数返回的数量(丢失记录的总数)相匹配,但通常情况并非如此 - 例如,我收到一条 Slack 消息但 check = 2。重新运行它会提供check = 1

有什么想法吗?

【问题讨论】:

    标签: scala apache-spark accumulator


    【解决方案1】:

    Spark 可以针对不同工作人员的相同数据多次运行一个方法,这意味着您计算每次成功 * 数据在任何工作人员上处理的次数。因此,对于相同数据的不同传递,您可以在累加器中获得不同的结果。

    在这种情况下,您不能使用累加器来获得准确的计数。对不起。 :(

    【讨论】:

    • 那为什么我只收到一条 Slack 消息?如果它被处理两次,那么我应该有两条消息
    • 嗯,抱歉,当时还不确定。虽然我没有过多地使用数据帧,但 iter 应该是一个分区上的所有数据,而不仅仅是 1 条记录。你确定你的成功只能是1或0吗?除此之外,我一无所有
    • 我认为这是不对的。因为他在 foreach 中运行它,这是一个动作,spark 保证累加器将准确更新一次,因此估值器应该是正确的。累加器仅在阶段完成时报告,因此在动作中运行时,即使需要重新运行,部分结果也不会影响最终值。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-06-19
    • 2019-01-02
    • 1970-01-01
    • 2011-01-22
    • 1970-01-01
    相关资源
    最近更新 更多