最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Python PySpark 核心實操入門案例指南

 更新時間:2026年04月03日 09:14:18   作者:別搶我的鍋包肉  
本文介紹了PySpark的核心機制,包括RDD的基本概念與創(chuàng)建方式、類型轉(zhuǎn)換與核心算子詳解、修改RDD分區(qū)的方法,以及解決MSVCR120.dll缺失問題的方案,感興趣的朋友跟隨小編一起看看吧

本文目標(biāo)

本文以 真實案例 + 深度拆解 的方式,帶你從零開始掌握 PySpark 的核心機制,涵蓋:

  • RDD 基本概念與創(chuàng)建
  • 數(shù)據(jù)類型轉(zhuǎn)換:T → TT → U
  • 核心算子詳解(map, flatMap, filter, distinct, sortBy, reduceByKey
  • 如何修改 RDD 分區(qū)?
  • 如何修復(fù) MSVCR120.dll 缺失?(Windows 專屬)
  • 如何解決 Spark 超時問題?
  • 兩個實戰(zhàn)案例:城市銷量排名 & 詞頻統(tǒng)計

一、PySpark 環(huán)境搭建基礎(chǔ)配置(關(guān)鍵?。?/h2>

特此強調(diào):Windows 用戶務(wù)必設(shè)置以下環(huán)境變量,否則幾乎所有操作都可能失??!

import os
# 1. 指定 Python 解釋器路徑(必須!)
os.environ["PYSPARK_PYTHON"] = "D:\\APP\\Anaconda\\envs\\spark_env\\python.exe"
# 2. Windows 超時問題(防止 task 超時崩潰)
os.environ['PYSPARK_TIMEOUT'] = '600'
os.environ['PYSPARK_DRIVER_TIMEOUT'] = '600'
# 3.  Hadoop native dll 加載導(dǎo)致 MSVCR120.dll 報錯
os.environ['PATH'] += os.pathsep + 'E:\\APP\\hadoop-3.4.2\\bin'
# 4. 強制設(shè)置 driver python(可選但推薦)
os.environ["PYSPARK_DRIVER_PYTHON"] = "D:\\APP\\Anaconda\\envs\\spark_env\\python.exe"

為什么這步如此重要?

  • MSVCR120.dllVisual C++ 2013 運行時庫,常在 Hadoop 的 hadoop.dll 調(diào)用失敗時出現(xiàn)。
  • 即使你不用 saveAsTextFile(),如果 PATH 中包含 hadoop/bin,Spark 仍會觸發(fā) native IO!
  • 終極建議:僅當(dāng)必須使用原生 IO 時才保留 PATH;否則刪除該路徑!

如何檢查到是 MSVCR120.dll 的問題: 這是我的報錯

from pyspark import SparkConf,SparkContext
import os
os.environ["PYSPARK_PYTHON"] = "D:\APP\Anaconda\envs\spark_env\python.exe"
os.environ['PATH'] += os.pathsep + 'E:\\APP\\hadoop-3.4.2\\bin'
conf = SparkConf().setMaster("local").setAppName("test_spark_app")
sc = SparkContext(conf=conf)
rdd = sc.parallelize([1,2,3,4,5])
rdd.saveAsTextFile("D:/output1")
sc.stop()D:\APP\Anaconda\envs\spark_env\python.exe D:\PythonProjects\python\day11_PySpark\數(shù)據(jù)輸出\輸出為文本文檔.py 
Setting default log level to "WARN".
To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel).
26/04/01 22:22:55 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
26/04/01 22:22:58 ERROR SparkHadoopWriter: Aborting job job_202604012222567528972640969878075_0003.
java.lang.UnsatisfiedLinkError: org.apache.hadoop.io.nativeio.NativeIO$Windows.access0(Ljava/lang/String;I)Z
	at org.apache.hadoop.io.nativeio.NativeIO$Windows.access0(Native Method)
	at org.apache.hadoop.io.nativeio.NativeIO$Windows.access(NativeIO.java:793)
	at org.apache.hadoop.fs.FileUtil.canRead(FileUtil.java:1249)
	at org.apache.hadoop.fs.FileUtil.list(FileUtil.java:1454)
	at org.apache.hadoop.fs.RawLocalFileSystem.listStatus(RawLocalFileSystem.java:601)
	at org.apache.hadoop.fs.FileSystem.listStatus(FileSystem.java:1972)
	at org.apache.hadoop.fs.FileSystem.listStatus(FileSystem.java:2014)
	at org.apache.hadoop.fs.ChecksumFileSystem.listStatus(ChecksumFileSystem.java:761)
	at org.apache.hadoop.fs.FileSystem.listStatus(FileSystem.java:1972)
	at org.apache.hadoop.fs.FileSystem.listStatus(FileSystem.java:2014)
	at org.apache.hadoop.mapreduce.lib.output.FileOutputCommitter.getAllCommittedTaskPaths(FileOutputCommitter.java:334)
	at org.apache.hadoop.mapreduce.lib.output.FileOutputCommitter.commitJobInternal(FileOutputCommitter.java:404)
	at org.apache.hadoop.mapreduce.lib.output.FileOutputCommitter.commitJob(FileOutputCommitter.java:377)
	at org.apache.hadoop.mapred.FileOutputCommitter.commitJob(FileOutputCommitter.java:136)
	at org.apache.hadoop.mapred.OutputCommitter.commitJob(OutputCommitter.java:291)
	at org.apache.spark.internal.io.HadoopMapReduceCommitProtocol.commitJob(HadoopMapReduceCommitProtocol.scala:192)
	at org.apache.spark.internal.io.SparkHadoopWriter$.$anonfun$write$3(SparkHadoopWriter.scala:100)
	at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
	at org.apache.spark.util.Utils$.timeTakenMs(Utils.scala:552)
	at org.apache.spark.internal.io.SparkHadoopWriter$.write(SparkHadoopWriter.scala:100)
	at org.apache.spark.rdd.PairRDDFunctions.$anonfun$saveAsHadoopDataset$1(PairRDDFunctions.scala:1091)
	at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
	at org.apache.spark.rdd.RDD.withScope(RDD.scala:407)
	at org.apache.spark.rdd.PairRDDFunctions.saveAsHadoopDataset(PairRDDFunctions.scala:1089)
	at org.apache.spark.rdd.PairRDDFunctions.$anonfun$saveAsHadoopFile$4(PairRDDFunctions.scala:1062)
	at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
	at org.apache.spark.rdd.RDD.withScope(RDD.scala:407)
	at org.apache.spark.rdd.PairRDDFunctions.saveAsHadoopFile(PairRDDFunctions.scala:1027)
	at org.apache.spark.rdd.PairRDDFunctions.$anonfun$saveAsHadoopFile$3(PairRDDFunctions.scala:1009)
	at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
	at org.apache.spark.rdd.RDD.withScope(RDD.scala:407)
	at org.apache.spark.rdd.PairRDDFunctions.saveAsHadoopFile(PairRDDFunctions.scala:1008)
	at org.apache.spark.rdd.PairRDDFunctions.$anonfun$saveAsHadoopFile$2(PairRDDFunctions.scala:965)
	at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
	at org.apache.spark.rdd.RDD.withScope(RDD.scala:407)
	at org.apache.spark.rdd.PairRDDFunctions.saveAsHadoopFile(PairRDDFunctions.scala:963)
	at org.apache.spark.rdd.RDD.$anonfun$saveAsTextFile$2(RDD.scala:1620)
	at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
	at org.apache.spark.rdd.RDD.withScope(RDD.scala:407)
	at org.apache.spark.rdd.RDD.saveAsTextFile(RDD.scala:1620)
	at org.apache.spark.rdd.RDD.$anonfun$saveAsTextFile$1(RDD.scala:1606)
	at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
	at org.apache.spark.rdd.RDD.withScope(RDD.scala:407)
	at org.apache.spark.rdd.RDD.saveAsTextFile(RDD.scala:1606)
	at org.apache.spark.api.java.JavaRDDLike.saveAsTextFile(JavaRDDLike.scala:564)
	at org.apache.spark.api.java.JavaRDDLike.saveAsTextFile$(JavaRDDLike.scala:563)
	at org.apache.spark.api.java.AbstractJavaRDDLike.saveAsTextFile(JavaRDDLike.scala:45)
	at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
	at sun.reflect.NativeMethodAccessorImpl.invoke(Unknown Source)
	at sun.reflect.DelegatingMethodAccessorImpl.invoke(Unknown Source)
	at java.lang.reflect.Method.invoke(Unknown Source)
	at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
	at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:374)
	at py4j.Gateway.invoke(Gateway.java:282)
	at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
	at py4j.commands.CallCommand.execute(CallCommand.java:79)
	at py4j.ClientServerConnection.waitForCommands(ClientServerConnection.java:182)
	at py4j.ClientServerConnection.run(ClientServerConnection.java:106)
	at java.lang.Thread.run(Unknown Source)
Traceback (most recent call last):
  File "D:\PythonProjects\python\day11_PySpark\數(shù)據(jù)輸出\輸出為文本文檔.py", line 9, in <module>
    rdd.saveAsTextFile("D:/output1")
  File "D:\APP\Anaconda\envs\spark_env\lib\site-packages\pyspark\rdd.py", line 3425, in saveAsTextFile
    keyed._jrdd.map(self.ctx._jvm.BytesToString()).saveAsTextFile(path)
  File "D:\APP\Anaconda\envs\spark_env\lib\site-packages\py4j\java_gateway.py", line 1322, in __call__
    return_value = get_return_value(
  File "D:\APP\Anaconda\envs\spark_env\lib\site-packages\py4j\protocol.py", line 326, in get_return_value
    raise Py4JJavaError(
py4j.protocol.Py4JJavaError: An error occurred while calling o33.saveAsTextFile.
: org.apache.spark.SparkException: Job aborted.
	at org.apache.spark.internal.io.SparkHadoopWriter$.write(SparkHadoopWriter.scala:106)
	at org.apache.spark.rdd.PairRDDFunctions.$anonfun$saveAsHadoopDataset$1(PairRDDFunctions.scala:1091)
	at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
	at org.apache.spark.rdd.RDD.withScope(RDD.scala:407)
	at org.apache.spark.rdd.PairRDDFunctions.saveAsHadoopDataset(PairRDDFunctions.scala:1089)
	at org.apache.spark.rdd.PairRDDFunctions.$anonfun$saveAsHadoopFile$4(PairRDDFunctions.scala:1062)
	at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
	at org.apache.spark.rdd.RDD.withScope(RDD.scala:407)
	at org.apache.spark.rdd.PairRDDFunctions.saveAsHadoopFile(PairRDDFunctions.scala:1027)
	at org.apache.spark.rdd.PairRDDFunctions.$anonfun$saveAsHadoopFile$3(PairRDDFunctions.scala:1009)
	at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
	at org.apache.spark.rdd.RDD.withScope(RDD.scala:407)
	at org.apache.spark.rdd.PairRDDFunctions.saveAsHadoopFile(PairRDDFunctions.scala:1008)
	at org.apache.spark.rdd.PairRDDFunctions.$anonfun$saveAsHadoopFile$2(PairRDDFunctions.scala:965)
	at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
	at org.apache.spark.rdd.RDD.withScope(RDD.scala:407)
	at org.apache.spark.rdd.PairRDDFunctions.saveAsHadoopFile(PairRDDFunctions.scala:963)
	at org.apache.spark.rdd.RDD.$anonfun$saveAsTextFile$2(RDD.scala:1620)
	at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
	at org.apache.spark.rdd.RDD.withScope(RDD.scala:407)
	at org.apache.spark.rdd.RDD.saveAsTextFile(RDD.scala:1620)
	at org.apache.spark.rdd.RDD.$anonfun$saveAsTextFile$1(RDD.scala:1606)
	at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
	at org.apache.spark.rdd.RDD.withScope(RDD.scala:407)
	at org.apache.spark.rdd.RDD.saveAsTextFile(RDD.scala:1606)
	at org.apache.spark.api.java.JavaRDDLike.saveAsTextFile(JavaRDDLike.scala:564)
	at org.apache.spark.api.java.JavaRDDLike.saveAsTextFile$(JavaRDDLike.scala:563)
	at org.apache.spark.api.java.AbstractJavaRDDLike.saveAsTextFile(JavaRDDLike.scala:45)
	at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
	at sun.reflect.NativeMethodAccessorImpl.invoke(Unknown Source)
	at sun.reflect.DelegatingMethodAccessorImpl.invoke(Unknown Source)
	at java.lang.reflect.Method.invoke(Unknown Source)
	at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
	at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:374)
	at py4j.Gateway.invoke(Gateway.java:282)
	at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
	at py4j.commands.CallCommand.execute(CallCommand.java:79)
	at py4j.ClientServerConnection.waitForCommands(ClientServerConnection.java:182)
	at py4j.ClientServerConnection.run(ClientServerConnection.java:106)
	at java.lang.Thread.run(Unknown Source)
Caused by: java.lang.UnsatisfiedLinkError: org.apache.hadoop.io.nativeio.NativeIO$Windows.access0(Ljava/lang/String;I)Z
	at org.apache.hadoop.io.nativeio.NativeIO$Windows.access0(Native Method)
	at org.apache.hadoop.io.nativeio.NativeIO$Windows.access(NativeIO.java:793)
	at org.apache.hadoop.fs.FileUtil.canRead(FileUtil.java:1249)
	at org.apache.hadoop.fs.FileUtil.list(FileUtil.java:1454)
	at org.apache.hadoop.fs.RawLocalFileSystem.listStatus(RawLocalFileSystem.java:601)
	at org.apache.hadoop.fs.FileSystem.listStatus(FileSystem.java:1972)
	at org.apache.hadoop.fs.FileSystem.listStatus(FileSystem.java:2014)
	at org.apache.hadoop.fs.ChecksumFileSystem.listStatus(ChecksumFileSystem.java:761)
	at org.apache.hadoop.fs.FileSystem.listStatus(FileSystem.java:1972)
	at org.apache.hadoop.fs.FileSystem.listStatus(FileSystem.java:2014)
	at org.apache.hadoop.mapreduce.lib.output.FileOutputCommitter.getAllCommittedTaskPaths(FileOutputCommitter.java:334)
	at org.apache.hadoop.mapreduce.lib.output.FileOutputCommitter.commitJobInternal(FileOutputCommitter.java:404)
	at org.apache.hadoop.mapreduce.lib.output.FileOutputCommitter.commitJob(FileOutputCommitter.java:377)
	at org.apache.hadoop.mapred.FileOutputCommitter.commitJob(FileOutputCommitter.java:136)
	at org.apache.hadoop.mapred.OutputCommitter.commitJob(OutputCommitter.java:291)
	at org.apache.spark.internal.io.HadoopMapReduceCommitProtocol.commitJob(HadoopMapReduceCommitProtocol.scala:192)
	at org.apache.spark.internal.io.SparkHadoopWriter$.$anonfun$write$3(SparkHadoopWriter.scala:100)
	at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
	at org.apache.spark.util.Utils$.timeTakenMs(Utils.scala:552)
	at org.apache.spark.internal.io.SparkHadoopWriter$.write(SparkHadoopWriter.scala:100)
	... 51 more
進程已結(jié)束,退出代碼為 1

由于找不到 MSVCR120.dll,無法繼續(xù)執(zhí)行代碼。重新安裝程序可能會解決此問題。 是 Windows 系統(tǒng)缺少 C++ 運行時庫 的典型問題。

?? 根本原因:
MSVCR120.dll 是 Microsoft Visual C++ 2013 Redistributable 的一部分。

它被 Hadoop 的原生 .dll 文件(如 hadoop.dll)所依賴,當(dāng)你使用 saveAsTextFile() 時,Spark 會嘗試加載 Hadoop 的本地庫,而這些庫又依賴于 VC++ 運行時。

即使你現(xiàn)在不用 saveAsTextFile(),你之前配置了 PATH += E:\APP\hadoop-3.4.2\bin,說明你還在用 Hadoop 二進制文件,所以問題依然存在。

我一開始懷疑是缺少winutils 環(huán)境配置 winutils配置可以看我上個文章 或者是 winutils.exe 權(quán)限不對
你的 winutils.exe 必須放在hadoop的bin目錄下 E:\APP\hadoop-3.4.2\bin\winutils.exe
然后檢查 權(quán)限
以 管理員身份 打開 CMD
一定要右鍵 → 以管理員身份運行命令提示符
. 執(zhí)行你這條命令

```bash
winutils.exe chmod 777 C:\tmp\Hive

可以用聯(lián)想的電腦異常修復(fù)直接修復(fù)

?? 二、RDD 基礎(chǔ)概念與創(chuàng)建方式 ??

什么是 RDD?

RDD(Resilient Distributed Dataset):彈性分布式數(shù)據(jù)集,是 Spark 的核心抽象。它是一個不可變、可并行、容錯的數(shù)據(jù)集合。

創(chuàng)建 RDD 的 6 種方式:

類型示例說明
Python 列表sc.parallelize([1,2,3])最常用
字典sc.parallelize({"k":1})返回 [{"k":1}]
元組sc.parallelize(("a","b"))返回 [("a","b")]
集合sc.parallelize([1,1,2,3])包含重復(fù)項
字符串sc.parallelize("abc")每個字符為一個元素
文件sc.textFile("path.txt")逐行讀取文本

推薦使用 parallelize() + textFile() 搭配處理本地數(shù)據(jù)。

三、核心算子詳解:類型轉(zhuǎn)換與功能解析

1.map(T)→U:傳入類型與返回類型不一致

作用:對每個元素進行函數(shù)映射,返回新值

rdd = sc.parallelize([1,2,3,4,5])
rdd2 = rdd.map(lambda x: x * 2)
print(rdd2.collect())  # 輸出: [2, 4, 6, 8, 10]
特點說明
T → U輸入是 int,輸出是 int(類型一致),但語義上是“變換”
按元素處理每個元素獨立轉(zhuǎn)換
無聚合不改變數(shù)據(jù)總量

適用場景:數(shù)據(jù)清洗、類型轉(zhuǎn)換、數(shù)學(xué)計算

2.map(T)→T:傳入類型與返回類型一致

作用:返回原類型,但可以修改內(nèi)容

rdd = sc.parallelize(["apple", "banana"])
rdd2 = rdd.map(lambda x: x.upper())
print(rdd2.collect())  # ['APPLE', 'BANANA']

雖然返回類型仍是 str,但內(nèi)容被修改 → 仍屬于 T → T

適用場景:字符串大小寫轉(zhuǎn)換、去除空格、文本標(biāo)準(zhǔn)化

3.flatMap(T)→U:扁平化處理,去嵌套

作用:將結(jié)果展開成平面列表

rdd = sc.parallelize(["hello world", "py spark"])
rdd2 = rdd.flatMap(lambda x: x.split(" "))
print(rdd2.collect())  # ['hello', 'world', 'py', 'spark']
區(qū)別map vs flatMap
map 返回 [[x], [y]]返回嵌套列表
flatMap 返回 [x, y]展開為單層列表

適用場景:文本分詞、多行拆分、列表合并

4.filter(T)→T:條件過濾

作用:根據(jù)返回 True/False 過濾元素

rdd = sc.parallelize([1,2,3,4,5,6])
rdd2 = rdd.filter(lambda x: x % 2 == 0)
print(rdd2.collect())  # [2, 4, 6]

適用場景:篩選城市、過濾商品類別、去除空值

5.distinct():去重

作用:返回去重后的集合(全局去重)

rdd = sc.parallelize([1,2,2,3,4,4,5])
rdd2 = rdd.distinct()
print(rdd2.collect())  # [1, 2, 3, 4, 5]

注意:distinct() 耗時高,會 Shuffle 所有數(shù)據(jù)!慎用于大數(shù)據(jù)集!

6.sortBy(keyfunc, ascending):排序

作用:按指定規(guī)則排序

rdd = sc.parallelize([1,5,3,2,4])
rdd2 = rdd.sortBy(lambda x: x, ascending=False)
print(rdd2.collect())  # [5, 4, 3, 2, 1]

適用場景:銷售排名、分?jǐn)?shù)排序、Top N 推薦

7.reduceByKey(func):分組聚合(鍵值對專屬)

前提:RDD 數(shù)據(jù)必須是 (K, V) 元組形式

rdd = sc.parallelize([("a",1),("a",2),("b",3),("b",4)])
rdd2 = rdd.reduceByKey(lambda a, b: a + b)
print(rdd2.collect())  # [('a', 3), ('b', 7)]
要求說明
必須是 (K, V)否則報錯
func(V,V) → V兩個 value 聚合出一個 value
本地預(yù)聚合減少網(wǎng)絡(luò)傳輸,效率極高

適用場景:統(tǒng)計商品銷量、城市銷售額、詞頻統(tǒng)計

四、修改 RDD 分區(qū)的 3 種方法

方法說明代碼示例
set("spark.default.parallelism", "1")設(shè)置全局并行度conf.set("spark.default.parallelism", "1")
numSlices=1parallelize 時指定分區(qū)數(shù)sc.parallelize([1,2,3], 1)
repartition(n)重新分區(qū)(常用于調(diào)優(yōu))rdd.repartition(2)

建議:小數(shù)據(jù)集用 numSlices=1 避免創(chuàng)建過多任務(wù)。

五、重難點攻克:如何修復(fù)MSVCR120.dll缺失問題???

報錯信息:

由于找不到 MSVCR120.dll,無法繼續(xù)執(zhí)行代碼

?? 根本原因:

  • Spark 調(diào)用 hadoop.dll 會依賴 VC++ 2013 運行時。
  • PATH 中包含 E:\APP\hadoop-3.4.2\bin → 啟動 native IO → 加載 DLL → 找不到函數(shù) → 報錯。

終極解決方案:

修復(fù) MSVCR120.dll 可以用聯(lián)想的電腦異常修復(fù)直接修復(fù)

六、案例實戰(zhàn):綜合數(shù)據(jù)分析

案例 1:城市銷量排名(P09_案例2.py)

需求:

  1. 讀取 sales.txt 文件(JSON 格式)
  2. 統(tǒng)計每個城市銷售額從大到小排序
  3. 獲取全部城市商品類別(去重)
  4. 查詢“北京”有哪些商品類別(過濾 + 去重)

步驟詳解:

import json
from pyspark import SparkContext, SparkConf
import os
os.environ["PYSPARK_PYTHON"] = "D:\\APP\\Anaconda\\envs\\spark_env\\python.exe"
os.environ['PYSPARK_TIMEOUT'] = '600'
os.environ['PYSPARK_DRIVER_TIMEOUT'] = '600'
conf = SparkConf().setMaster("local[*]").setAppName("sales_analysis")
sc = SparkContext(conf=conf)
# 1. 讀取文件并轉(zhuǎn)為字典
file_rdd = sc.textFile("D:\\python_testdata\\sales.txt")
dict_rdd = file_rdd.map(lambda x: json.loads(x.rstrip().rstrip(',')))
# 2. 各城市銷售額排名
city_money = dict_rdd.map(lambda x: (x['areaName'], int(x['money'])))
city_total = city_money.reduceByKey(lambda a, b: a + b)
city_rank = city_total.sortBy(lambda x: x[1], ascending=False)
print("各城市銷售額排名:")
print(city_rank.collect())
# 3. 全部城市有哪些商品類別(去重)
categories = dict_rdd.map(lambda x: x['category'])
unique_categories = categories.distinct()
print("全部商品類別:")
print(unique_categories.collect())
# 4. 北京市的商品類別
beijing_rdd = dict_rdd.filter(lambda x: x['areaName'] == '北京')
beijing_cats = beijing_rdd.map(lambda x: x['category']).distinct()
print("北京市售賣的商品類別:")
print(beijing_cats.collect())
sc.stop()

知識點總結(jié):

  • json.loads():解析 JSON 字符串
  • map:提取字段
  • reduceByKey:聚合銷售額
  • sortBy:排序
  • filter + distinct:組合過濾去重

案例 2:詞頻統(tǒng)計(P05_PySpark案例.py)

需求:統(tǒng)計words.txt中每個詞出現(xiàn)次數(shù)

from pyspark import SparkContext, SparkConf
import os
os.environ["PYSPARK_PYTHON"] = "D:\\APP\\Anaconda\\envs\\spark_env\\python.exe"
os.environ['PYSPARK_TIMEOUT'] = '600'
conf = SparkConf().setMaster("local[*]").setAppName("word_count")
sc = SparkContext(conf=conf)
rdd = sc.textFile("D:\\python_testdata\\words.txt")
words = rdd.flatMap(lambda x: x.split(" "))
word_pairs = words.map(lambda x: (x, 1))
word_count = word_pairs.reduceByKey(lambda a, b: a + b)
print("詞頻統(tǒng)計結(jié)果:")
print(word_count.collect())
sc.stop()

知識點:

  • flatMap 分詞
  • map(word, 1)
  • reduceByKey 聚合
  • 典型的 MapReduce 模型

七、總結(jié)與建議

項目推薦做法
數(shù)據(jù)輸出? 使用 collect() + open() 寫文件,避免 saveAsTextFile()
環(huán)境配置? 設(shè)置 PYSPARK_TIMEOUT + spark.hadoop.io.native.io.enabled=false
修復(fù) DLL 錯誤? 刪除 PATH 中 hadoop/bin,不加載 native IO
超時問題? 設(shè)置 PYSPARK_TIMEOUT=600,防止任務(wù)中斷
小數(shù)據(jù)集? 用 numSlices=1 控制分區(qū)
大數(shù)據(jù)集? 用 repartition() 調(diào)優(yōu)

到此這篇關(guān)于Python PySpark 核心實操入門案例指南的文章就介紹到這了,更多相關(guān)Python PySpark 入門內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

最新評論

商城县| 新平| 天祝| 汶川县| 左权县| 泾川县| 织金县| 杭州市| 花莲县| 阿坝县| 上蔡县| 独山县| 湘阴县| 当涂县| 门源| 巴南区| 诸暨市| 高清| 西藏| 益阳市| 黄冈市| 肥东县| 左贡县| 兴和县| 陕西省| 永善县| 新绛县| 汉沽区| 年辖:市辖区| 资兴市| 凉山| 岳池县| 石河子市| 长寿区| 海原县| 铁岭市| 甘肃省| 太谷县| 新巴尔虎右旗| 曲麻莱县| 德保县|