我正在嘗試計算檔案中存在的有效和無效資料的數量。下面是執行相同操作的代碼,
val badDataCountAcc = spark.sparkContext.longAccumulator("BadDataAcc")
val goodDataCountAcc = spark.sparkContext.longAccumulator("GoodDataAcc")
val dataframe = spark
.read
.format("csv")
.option("header", true)
.option("inferSchema", true)
.load(path)
.filter(data => {
val matcher = regex.matcher(data.toString())
if (matcher.find()) {
goodDataCountAcc.add(1)
println("GoodDataCountAcc: " goodDataCountAcc.value)
true
} else {
badDataCountAcc.add(1)
println("BadDataCountAcc: " badDataCountAcc.value)
false
}
}
)
.withColumn("FileName", input_file_name())
dataframe.show()
val filename = dataframe
.select("FileName")
.distinct()
val name = filename.collectAsList().get(0).toString()
println("" filename)
println("Bad data Count Acc: " badDataCountAcc.value)
println("Good data Count Acc: " goodDataCountAcc.value)
我為具有 2 個有效資料和 3 個無效資料的示例資料運行了此代碼。在我列印計數的過濾器內部,值是正確的。但是當我列印計數值時,在過濾器之外,好的資料是 4,壞資料是 6。
問題:
- 當我在最后洗掉 withColumn 陳述句時 - 連同計算不同檔案名的代碼 - 值被正確列印。我不確定為什么?
- 我也確實需要獲取輸入檔案名。在這里最好的方法是什么?
uj5u.com熱心網友回復:
首先,Accumulator 屬于 RDD API,而您使用的是 Dataframes。資料幀最終被編譯為 RDD,但它們處于更高的抽象級別。在這種情況下,最好使用聚合而不是累加器。
來自Spark Accumulators 檔案:
對于僅在動作內部執行的累加器更新,Spark 保證每個任務對累加器的更新只會應用一次,即重新啟動的任務不會更新值。在轉換中,用戶應注意,如果重新執行任務或作業階段,每個任務的更新可能會應用多次。
累加器不會改變 Spark 的惰性求值模型。如果它們在對 RDD 的操作中被更新,則它們的值僅在該 RDD 被計算為操作的一部分時才會更新。因此,當在 map() 等惰性轉換中進行累加器更新時,不能保證執行累加器更新。下面的代碼片段演示了這個屬性:
您的 DataFramefilter將遵守 RDD filter,這不是一個動作,而是一個轉換(因此是懶惰的),因此這種僅一次的保證不適用于您的情況。你的代碼執行多少次取決于實作,并且可能會隨著 Spark 版本而改變,所以你不應該依賴它。
關于你的兩個問題:
(編輯前)無法根據您的代碼段回答此問題,因為它不包含任何操作。它甚至是您使用的確切代碼片段嗎?我懷疑如果您實際上執行了您發布的代碼,除了缺少的匯入之外沒有任何添加,它應該列印 0 兩次,因為沒有執行任何操作。無論哪種方式,您都應該始終假設 RDD 轉換中的累加器可能會執行多次(如果它在可能被優化的 DataFrame 操作中,甚至根本不會執行)。
您的使用
withColumn方法非常好。
我建議使用 DataFrame 運算式和聚合(如果您愿意,也可以使用等效的 Spark SQL)。可以使用rlike, 使用列而不是依賴來完成正則運算式匹配toString(),例如.withColumn("IsGoodData", $"myColumn1".rlike(regex1) && $"myColumn2".rlike(regex2))。
然后你可以使用聚合來計算好記錄和壞記錄dataframe.groupBy($"IsGoodData").count()
編輯:通過附加行,您的第一個問題的答案也很清楚:第一次來自dataframe.show(),第二次來自filename.collectAsList(),您可能也將其洗掉,因為它取決于添加的列。請確保您了解 Spark 轉換和操作與 Spark 的惰性評估模型之間的區別。否則你不會對它很滿意:-)
轉載請註明出處,本文鏈接:https://www.uj5u.com/yidong/496940.html
上一篇:來自按時間戳排序的兩個不同kafka主題的Spark聚合事件
下一篇:型別引數的功能障礙
