欢迎来到尧图网

客户服务 关于我们

您的位置:首页 > 财经 > 创投人物 > Windows环境下开发pyspark程序

Windows环境下开发pyspark程序

2025/4/18 20:43:13 来源:https://blog.csdn.net/W_chuanqi/article/details/147004318  浏览:    关键词:Windows环境下开发pyspark程序

Windows环境下开发pyspark程序

一、环境准备

1.1. Anaconda/Miniconda(Python环境)

如果不怕包的版本管理混乱,可以直接使用已有的Python环境。

需要安装anaconda/miniconda(python3.8版本以上):Anaconda下载安装及老版本选择(超详细)
使用conda新建一个虚拟环境用于PySpark开发:Python虚拟环境(windows)

首先,我们新建一个 pyspark_env 文件夹,作为虚拟环境的存放路径(也可以不用,conda创建虚拟环境时检测到没有会自动新建):
在这里插入图片描述
创建环境并指定路径:

conda create -p E:\penv\pyspark_env python=3.9

在这里插入图片描述
创建完成:
在这里插入图片描述
激活环境:

conda activate E:\penv\pyspark_env

在这里插入图片描述
安装pyspark:

pip install pyspark -i https://pypi.tuna.tsinghua.edu.cn/simple/

在这里插入图片描述
安装psutil

pip install psutil

在这里插入图片描述

1.2. JDK

请注意,PySpark需要Java 8(不包括8u371之前版本)、11或17,并且JAVA_HOME需要正确设置。设置JAVA安装路径的时候不要有空格,否则会报错。
在这里插入图片描述
参考这篇文章:JDK8卸载与安装教程(超详细)

1.3. 安装hadoop

(1)下载
进入hadoop安装包下载地址,这里选择的是hadoop-3.3.6.tar.gz版本:
在这里插入图片描述
(2)解压
对下载好的文件进行解压,将其解压放在个人想存放的目录中(记住路径,以便配置环境变量)。
在这里插入图片描述
在这里插入图片描述
解压成功:
在这里插入图片描述
(3)配置环境变量

HADOOP_HOME

在这里插入图片描述

%HADOOP_HOME%\bin

在这里插入图片描述
此时bin目录( E:\hadoop-3.3.6\bin)下没有 hadoop.dll及winutils.exe文件:
在这里插入图片描述

  • 需要进行下载winutils :https://soft.3dmgame.com/down/204154.html
    在这里插入图片描述

  • 解压文件,选择hadoop版本对应的文件夹bin目录下的hadoop.dll和winutils.exe文件
    在这里插入图片描述

  • 将hadoop.dll和winutils.exe 拷贝到E:\hadoop-3.3.6\bin 、C:\Windows\System32下(两个文件各拷贝一份到两个目录中)
    在这里插入图片描述
    在这里插入图片描述
    (4)环境测试

二、新建一个Python项目

2.1. 创建项目并配置解释器

新建一个项目,项目名为pyspark
在这里插入图片描述
添加新的解释器(找到虚拟环境中的python.exe):
在这里插入图片描述
在这里插入图片描述
创建项目:
在这里插入图片描述

2.2. 创建目录文件

main :用于存放每天开发的一些代码文件
resources :用于存放程序中需要用到的配置文件
datas :用于存放每天用到的一些数据文件
test :用于存放测试时的一些代码文件

在这里插入图片描述
在这里插入图片描述
在这里插入图片描述

2.3. 环境测试

import os
from pyspark import SparkContext, SparkConf  # 导入pyspark模块if __name__ == '__main__':# 配置环境os.environ['JAVA_HOME'] = 'C:/Program Files/Java/jdk-1.8'# 配置Hadoop的路径,就是前面解压的那个路径os.environ['HADOOP_HOME'] = 'E:/hadoop-3.3.6'# 配置Python解析器的路径os.environ['PYSPARK_PYTHON'] = 'E:/penv/pyspark_env/python.exe' os.environ['PYSPARK_DRIVER_PYTHON'] = 'E:/penv/pyspark_env/python.exe'# 获取 conf 对象# setMaster  按照什么模式运行,local  bigdata01:7077  yarn#  local[2]  使用2核CPU   * 你本地资源有多少核就用多少核#  appName 任务的名字conf = SparkConf().setMaster("local[*]").setAppName("第一个Spark程序")# 假如我想设置压缩# conf.set("spark.eventLog.compression.codec","snappy")# 根据配置文件,得到一个SC对象,第一个conf 是 形参的名字,第二个conf 是实参的名字sc = SparkContext(conf=conf)print(sc)# 使用完后,记得关闭sc.stop()

输出结果:
在这里插入图片描述

三、WordCount案例

3.1 数据准备

这里我使用文心一言生成了一份数据,用来测试WordCount
数据如下所示:

Hello World! This is a simple WordCount example. The WordCount program is used to count the frequency of words in a given text.Let's analyze this example: "Hello World!" Hello again, World! Notice how the word 'Hello' appears multiple times, as does 'World'.The program should ignore case sensitivity, meaning 'Hello' and 'hello' should be treated as the same word. Additionally, punctuation marks like commas, periods, and exclamation points should not affect the word count.In summary, a WordCount program takes text as input and outputs a list of words along with their corresponding frequencies. For instance, the word 'Hello' might appear 3 times, while 'World' appears 2 times in this example.

数据特点

  • 重复单词:Hello 和 World 多次出现。
  • 标点符号:包含逗号、句号和感叹号等标点符号。
  • 大小写混合:Hello 和 hello 应被视为同一个单词。
  • 自然语言结构:包含简单句子和段落,模拟真实文本。

3.2 代码实现

代码实现如下所示:

import os
import re
from pyspark import SparkContext, SparkConf  # 导入pyspark模块if __name__ == '__main__':# 配置环境os.environ['JAVA_HOME'] = 'C:/Program Files/Java/jdk-1.8'# 配置Hadoop的路径,就是前面解压的那个路径os.environ['HADOOP_HOME'] = 'E:/hadoop-3.3.6'# 配置base环境Python解析器的路径os.environ['PYSPARK_PYTHON'] = 'E:/penv/pyspark_env/python.exe'  # 配置base环境Python解析器的路径os.environ['PYSPARK_DRIVER_PYTHON'] = 'E:/penv/pyspark_env/python.exe'# 获取 conf 对象# setMaster  按照什么模式运行,local  bigdata01:7077  yarn#  local[2]  使用2核CPU   * 你本地资源有多少核就用多少核#  appName 任务的名字conf = SparkConf().setMaster("local[*]").setAppName("WordCount")# 假如我想设置压缩# conf.set("spark.eventLog.compression.codec","snappy")# 根据配置文件,得到一个SC对象,第一个conf 是 形参的名字,第二个conf 是实参的名字sc = SparkContext(conf=conf)fileRdd = sc.textFile("../datas/wordcount/word.txt")  # 读取数据rsRdd = fileRdd \.filter(lambda line: len(line.strip()) > 0) \.flatMap(lambda line: line.strip().split(r" ")) \.map(lambda word: (word, 1)) \.reduceByKey(lambda a, b: a + b)rsRdd.saveAsTextFile("../output")# 使用完后,记得关闭sc.stop()

输出结果:
在这里插入图片描述
代码解析:

  • filter(lambda line: len(line.strip()) > 0):过滤掉空行;
    在这里插入图片描述
  • flatMap(lambda line:line.strip().split(r"")):将每一行多个单词转换为一行一个单词,r 的作用是告诉 Python 将字符串按原始字符串处理,避免转义字符的干扰。
    在这里插入图片描述
  • .map(lambda word: (word, 1)):将每个单词转换成KeyValue的二元组(word,1)
    在这里插入图片描述
  • reduceByKey(lambda a, b: a + b):先根据key值进行分组,然后再进行聚合。
    在这里插入图片描述

3.3 代码改进

虽然代码实现出来了简单的WordCount,但是没有达到我们想要的预期,主要有以下几点需要改进:

  • 单词前后的符号无法处理,导致一个单词分成了不同的组。
    在这里插入图片描述
  • 对单词的大小写不敏感,如:Hello和hello应视为一个词。
    在这里插入图片描述
3.3.1 解决标点符号

对于标点符号,我们可以使用正则表达式进行处理。
下面是正则表达式的一个测试用例:

import retext = "你好,世界!这是一个测试文本。"
# 使用正则表达式去除标点符号
result = re.sub(r'[^\w\s]', '', text)
print(result)  # 输出:你好世界这是一个测试文本

其中:

  • [^\w\s] 匹配所有非字母、数字和空格的字符(即标点符号)。
  • re.sub() 将匹配的字符替换为空字符串。
3.3.2 解决大小写字母

对单词的大小写不敏感,我们可以采取以下措施。

  • 全部字母大写或者小写:使用upper()或者lower()函数。
text = "Hello World"
upper_text = text.upper()
lower_text = text.lower()
print(upper_text)  # HELLO WORLD
print(lower_text)  # hello world
  • 首字母大写,其余字母小写:
  1. 使用 capitalize() 方法 capitalize() 方法会将字符串的第一个字符转换为大写,其余字符转换为小写。
text = "hello world"
capitalized_text = text.capitalize()
print(capitalized_text)  # 输出: Hello world
  1. 使用 title() 方法 如果你希望字符串中每个单词的首字母都大写,可以使用 title() 方法。
text = "hello world"
title_text = text.title()
print(title_text)  # 输出: Hello World
3.3.3 代码实现

这里我们采用正则表达式对标点符号进行处理,使用title()方法处理字母大小写。那么,改进后的代码如下:

import os
import re
from pyspark import SparkContext, SparkConf  # 导入pyspark模块if __name__ == '__main__':# 配置环境os.environ['JAVA_HOME'] = 'C:/Program Files/Java/jdk-1.8'# 配置Hadoop的路径,就是前面解压的那个路径os.environ['HADOOP_HOME'] = 'E:/hadoop-3.3.6'# 配置base环境Python解析器的路径os.environ['PYSPARK_PYTHON'] = 'E:/penv/pyspark_env/python.exe'  # 配置base环境Python解析器的路径os.environ['PYSPARK_DRIVER_PYTHON'] = 'E:/penv/pyspark_env/python.exe'# 获取 conf 对象# setMaster  按照什么模式运行,local  bigdata01:7077  yarn#  local[2]  使用2核CPU   * 你本地资源有多少核就用多少核#  appName 任务的名字conf = SparkConf().setMaster("local[*]").setAppName("WordCount")# 假如我想设置压缩# conf.set("spark.eventLog.compression.codec","snappy")# 根据配置文件,得到一个SC对象,第一个conf 是 形参的名字,第二个conf 是实参的名字sc = SparkContext(conf=conf)fileRdd = sc.textFile("../datas/wordcount/word.txt")  # 读取数据rsRdd = fileRdd \.filter(lambda line: len(line.strip()) > 0) \.flatMap(lambda line: re.sub(r'[^\w\s]', '', line.strip()).split()) \.map(lambda word: (word.title(), 1)) \.reduceByKey(lambda a, b: a + b)rsRdd.saveAsTextFile("../output3")# 使用完后,记得关闭sc.stop()

输出结果为:

('Hello', 7)
('World', 5)
('Wordcount', 3)
('Example', 3)
('The', 7)
('Program', 3)
('Count', 2)
('Of', 2)
('Words', 2)
('Lets', 1)
('Analyze', 1)
('Again', 1)
('Appears', 2)
('Should', 3)
('Ignore', 1)
('Sensitivity', 1)
('And', 3)
('Same', 1)
('Punctuation', 1)
('Marks', 1)
('Periods', 1)
('Not', 1)
('Affect', 1)
('Summary', 1)
('Takes', 1)
('Input', 1)
('Outputs', 1)
('List', 1)
('Corresponding', 1)
('Instance', 1)
('Might', 1)
('This', 3)
('Is', 2)
('A', 4)
('Simple', 1)
('Used', 1)
('To', 1)
('Frequency', 1)
('In', 3)
('Given', 1)
('Text', 2)
('Notice', 1)
('How', 1)
('Word', 4)
('Multiple', 1)
('Times', 3)
('As', 3)
('Does', 1)
('Case', 1)
('Meaning', 1)
('Be', 1)
('Treated', 1)
('Additionally', 1)
('Like', 1)
('Commas', 1)
('Exclamation', 1)
('Points', 1)
('Along', 1)
('With', 1)
('Their', 1)
('Frequencies', 1)
('For', 1)
('Appear', 1)
('3', 1)
('While', 1)
('2', 1)

四、数据去重案例

4.1 数据准备

这里提供了csv版本的数据:

ID ,  Name    ,  Email                 ,  Phone        ,  Address
1  ,  Alice   ,  alice@example.com     ,  123-456-7890 ,  123 Main St
2  ,  Bob     ,  bob@example.com       ,  234-567-8901 ,  456 Elm St
3  ,  Alice   ,  alice@example.com     ,  123-456-7890 ,  123 Main St
4  ,  Charlie ,  charlie@example.com   ,  345-678-9012 ,  789 Oak St
5  ,  David   ,  david@example.com     ,  456-789-0123 ,  101 Pine St
6  ,  Alice   ,  alice.new@example.com ,  123-456-7890 ,  123 Main St (new addr)
7  ,  Bob     ,  bob@example.com       ,  234-567-8901 ,  456 Elm St (alt addr)
8  ,  Eve     ,  eve@example.com       ,  567-890-1234 ,  202 Maple St
9  ,  Charlie ,  charlie@example.com   ,  345-678-9012 ,  789 Oak St

4.2 去重规则

  1. 完全匹配去重:如果两行数据的所有字段都相同,则认为是重复项,保留其中一行。
  2. 部分匹配去重(可选):如果某些字段(如 Name 和 Email)相同,但其他字段(如 Phone 和 Address)不同,可以根据业务需求决定是否视为重复项。

在此示例中,我们仅考虑完全匹配去重。

4.3 代码实现

方法一:使用PySpark中dataframe进行实现:

import os
from pyspark.sql import SparkSessionif __name__ == '__main__':# 配置环境os.environ['JAVA_HOME'] = 'C:/Program Files/Java/jdk-1.8'# 配置Hadoop的路径,就是前面解压的那个路径os.environ['HADOOP_HOME'] = 'E:/hadoop-3.3.6'# 配置base环境Python解析器的路径os.environ['PYSPARK_PYTHON'] = 'E:/penv/pyspark_env/python.exe'os.environ['PYSPARK_DRIVER_PYTHON'] = 'E:/penv/pyspark_env/python.exe'# 创建SparkSessionspark = SparkSession.builder \.appName("Data Deduplication") \.getOrCreate()# 读取CSV文件csv_file_path = "../datas/data deduplication/data.csv"# header=True表示第一行作为列名,inferSchema=True尝试自动推断数据类型。df = spark.read.csv(csv_file_path, header=True, inferSchema=True)    # 显示原始数据print("原始数据:")df.show()# 获取所有列名并排除ID字段columns_to_check = df.columns[1:]# 去除重复行(忽略ID字段)# 使用dropDuplicates()函数基于columns_to_check列表中的列名去除重复行。这意味着如果两行在这些列上的值完全相同,则只保留一行。df_no_duplicates = df.dropDuplicates(subset=columns_to_check)# 显示去重后的数据print("去重后的数据:")df_no_duplicates.show()# 如果需要保存去重后的数据到新的CSV文件output_csv_file_path = "../datas/data deduplication/deduplicated_data.csv"df_no_duplicates.write.csv(output_csv_file_path, header=True, mode="overwrite")# 停止SparkSessionspark.stop()

输出结果为:
在这里插入图片描述
在这里插入图片描述
方法二:使用PySaprk中的SQL进行实现。

import os
from pyspark.sql import SparkSessionif __name__ == '__main__':# 配置环境os.environ['JAVA_HOME'] = 'C:/Program Files/Java/jdk-1.8'# 配置Hadoop的路径,就是前面解压的那个路径os.environ['HADOOP_HOME'] = 'E:/hadoop-3.3.6'# 配置base环境Python解析器的路径os.environ['PYSPARK_PYTHON'] = 'E:/penv/pyspark_env/python.exe'os.environ['PYSPARK_DRIVER_PYTHON'] = 'E:/penv/pyspark_env/python.exe'# 创建SparkSessionspark = SparkSession.builder \.appName("Data Deduplication") \.getOrCreate()# 读取CSV文件csv_file_path = "../datas/data deduplication/data.csv"# header=True表示第一行作为列名,inferSchema=True尝试自动推断数据类型。df = spark.read.csv(csv_file_path, header=True, inferSchema=True)# 获取所有列名并排除ID字段columns_to_check = df.columns[1:]# 创建一个不包含ID字段的DataFramedf = df.select(columns_to_check)# 创建一个临时视图df.createOrReplaceTempView("my_table")spark.sql("select DISTINCT * from my_table").show()# 停止SparkSessionspark.stop()

输出结果:
在这里插入图片描述
但是这个没有对应的ID列。
方法三

import os
from pyspark.sql import SparkSessionif __name__ == '__main__':# 配置环境os.environ['JAVA_HOME'] = 'C:/Program Files/Java/jdk-1.8'# 配置Hadoop的路径,就是前面解压的那个路径os.environ['HADOOP_HOME'] = 'E:/hadoop-3.3.6'# 配置base环境Python解析器的路径os.environ['PYSPARK_PYTHON'] = 'E:/penv/pyspark_env/python.exe'os.environ['PYSPARK_DRIVER_PYTHON'] = 'E:/penv/pyspark_env/python.exe'# 创建SparkSessionspark = SparkSession.builder \.appName("Data Deduplication") \.getOrCreate()# 读取CSV文件csv_file_path = "../datas/data deduplication/data.csv"# header=True表示第一行作为列名,inferSchema=True尝试自动推断数据类型。df = spark.read.csv(csv_file_path, header=True, inferSchema=True)# 创建一个临时视图df.createOrReplaceTempView("tt")spark.sql("""SELECT * FROM ttWHERE ID  IN (select min(ID) from tt group by Name,Email,Phone,Address)""").show()# 停止SparkSessionspark.stop()

输出结果:
在这里插入图片描述

问题

1. 测试hadoop出现错误

在这里插入图片描述
原因分析:这时候,多半是因为你的java环境变量路径含有空格。
解决方法
(1)找到hadoop\etc\hadoop这个目录下的hadoop-env.cmd这个命令脚本。
在这里插入图片描述
然后,右键,编辑/notpad ++ ,进入编辑页面:
在这里插入图片描述
修改JAVA_HOME,我的JAVA的安装路径为:C:\Program Files\Java\jdk-1.8
在这里插入图片描述
添加引号:
在这里插入图片描述
查看hadoop版本:
在这里插入图片描述

2. Please install psutil

运行代码,出现下面的情况:

E:\penv\pyspark_env\Lib\site-packages\pyspark\python\lib\pyspark.zip\pyspark\shuffle.py:65: UserWarning: Please install psutil to have better support with spilling
E:\penv\pyspark_env\Lib\site-packages\pyspark\python\lib\pyspark.zip\pyspark\shuffle.py:65: UserWarning: Please install psutil to have better support with spilling
E:\penv\pyspark_env\Lib\site-packages\pyspark\python\lib\pyspark.zip\pyspark\shuffle.py:65: UserWarning: Please install psutil to have better support with spilling
E:\penv\pyspark_env\Lib\site-packages\pyspark\python\lib\pyspark.zip\pyspark\shuffle.py:65: UserWarning: Please install psutil to have better support with spilling进程已结束,退出代码0

在这里插入图片描述
解决方案,安装这个包:

pip install psutil

参考

  1. Hadoop高手之路4-HDFS
  2. Anaconda下载安装及老版本选择(超详细)
  3. Python虚拟环境(windows)
  4. JDK8卸载与安装教程(超详细)
  5. Windows环境本地配置pyspark环境详细教程
  6. win10下执行Hadoop命令报错:系统找不到指定的路径。Error: JAVA_HOME is incorrectly set. Please update D:
  7. sparkRDD编程实战

版权声明:

本网仅为发布的内容提供存储空间,不对发表、转载的内容提供任何形式的保证。凡本网注明“来源:XXX网络”的作品,均转载自其它媒体,著作权归作者所有,商业转载请联系作者获得授权,非商业转载请注明出处。

我们尊重并感谢每一位作者,均已注明文章来源和作者。如因作品内容、版权或其它问题,请及时与我们联系,联系邮箱:809451989@qq.com,投稿邮箱:809451989@qq.com

热搜词