【发布时间】: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