使用 PySpark 的 Big Data 基礎
Upendra Devisetty
Science Analyst, CyVerse
reduce(func) 動作用於彙總一般 RDD 的元素
此函式需同時具備交換律(運算元順序改變不影響結果)與結合律
在 PySpark 中 reduce() 的範例
x = [1,3,4,6]
RDD = sc.parallelize(x)
RDD.reduce(lambda x, y : x + y)
14
saveAsTextFile() 會將 RDD 以文字檔儲存在一個目錄中,每個分割區各成一檔RDD.saveAsTextFile("tempFile")
coalesce() 將 RDD 儲存為單一文字檔RDD.coalesce(1).saveAsTextFile("tempFile")
適用於 PySpark pair RDD 的 RDD 動作
pair RDD 動作會利用鍵值資料
常見的 pair RDD 動作包含
countByKey()
collectAsMap()
countByKey() 僅適用於型別 (K, V)
countByKey() 會計算每個鍵的元素數量
在簡單清單上的 countByKey() 範例
rdd = sc.parallelize([("a", 1), ("b", 1), ("a", 1)])
for kee, val in rdd.countByKey().items():
print(kee, val)
('a', 2)
('b', 1)
collectAsMap() 會將 RDD 的鍵值對回傳為字典
在簡單 tuple 上的 collectAsMap() 範例
sc.parallelize([(1, 2), (3, 4)]).collectAsMap()
{1: 2, 3: 4}
使用 PySpark 的 Big Data 基礎