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

Python PySpark 核心實(shí)操入門案例指南

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

本文目標(biāo)

本文以 真實(shí)案例 + 深度拆解 的方式,帶你從零開(kāi)始掌握 PySpark 的核心機(jī)制,涵蓋:

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

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

特此強(qiáng)調(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 超時(shí)問(wèn)題(防止 task 超時(shí)崩潰)
os.environ['PYSPARK_TIMEOUT'] = '600'
os.environ['PYSPARK_DRIVER_TIMEOUT'] = '600'
# 3.  Hadoop native dll 加載導(dǎo)致 MSVCR120.dll 報(bào)錯(cuò)
os.environ['PATH'] += os.pathsep + 'E:\\APP\\hadoop-3.4.2\\bin'
# 4. 強(qiáng)制設(shè)置 driver python(可選但推薦)
os.environ["PYSPARK_DRIVER_PYTHON"] = "D:\\APP\\Anaconda\\envs\\spark_env\\python.exe"

為什么這步如此重要?

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

如何檢查到是 MSVCR120.dll 的問(wèn)題: 這是我的報(bào)錯(cuò)

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
進(jìn)程已結(jié)束,退出代碼為 1

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

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

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

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

我一開(kāi)始懷疑是缺少winutils 環(huán)境配置 winutils配置可以看我上個(gè)文章 或者是 winutils.exe 權(quán)限不對(duì)
你的 winutils.exe 必須放在hadoop的bin目錄下 E:\APP\hadoop-3.4.2\bin\winutils.exe
然后檢查 權(quán)限
以 管理員身份 打開(kāi) CMD
一定要右鍵 → 以管理員身份運(yùn)行命令提示符
. 執(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 的核心抽象。它是一個(gè)不可變、可并行、容錯(cuò)的數(shù)據(jù)集合。

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

類型示例說(shuō)明
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ù)項(xiàng)
字符串sc.parallelize("abc")每個(gè)字符為一個(gè)元素
文件sc.textFile("path.txt")逐行讀取文本

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

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

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

作用:對(duì)每個(gè)元素進(jìn)行函數(shù)映射,返回新值

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

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

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

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

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

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

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]展開(kāi)為單層列表

適用場(chǎng)景:文本分詞、多行拆分、列表合并

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

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

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

適用場(chǎng)景:篩選城市、過(guò)濾商品類別、去除空值

5.distinct():去重

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

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

注意:distinct() 耗時(shí)高,會(huì) 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]

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

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

前提: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)]
要求說(shuō)明
必須是 (K, V)否則報(bào)錯(cuò)
func(V,V) → V兩個(gè) value 聚合出一個(gè) value
本地預(yù)聚合減少網(wǎng)絡(luò)傳輸,效率極高

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

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

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

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

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

報(bào)錯(cuò)信息:

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

?? 根本原因:

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

終極解決方案:

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

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

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

需求:

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

步驟詳解:

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()

知識(shí)點(diǎn)總結(jié):

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

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

需求:統(tǒng)計(jì)words.txt中每個(gè)詞出現(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)計(jì)結(jié)果:")
print(word_count.collect())
sc.stop()

知識(shí)點(diǎn):

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

七、總結(jié)與建議

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

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

相關(guān)文章

最新評(píng)論

江北区| 遂宁市| 克什克腾旗| 那曲县| 永修县| 绥阳县| 务川| 嘉祥县| 拉孜县| 天台县| 江永县| 册亨县| 新宾| 扎兰屯市| 广水市| 沭阳县| 鄂伦春自治旗| 巴彦淖尔市| 闻喜县| 东宁县| 晋城| 原平市| 台北市| 保康县| 从江县| 上蔡县| 法库县| 阿坝县| 东明县| 内黄县| 孙吴县| 祁连县| 广德县| 宣汉县| 武清区| 巨野县| 商洛市| 绥棱县| 资溪县| 巩留县| 新闻|