当前位置: 首页 > news >正文

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 re

text = "你好,世界!这是一个测试文本。"
# 使用正则表达式去除标点符号
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 SparkSession

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'
    os.environ['PYSPARK_DRIVER_PYTHON'] = 'E:/penv/pyspark_env/python.exe'

    # 创建SparkSession
    spark = 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")

    # 停止SparkSession
    spark.stop()

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

import os
from pyspark.sql import SparkSession

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'
    os.environ['PYSPARK_DRIVER_PYTHON'] = 'E:/penv/pyspark_env/python.exe'

    # 创建SparkSession
    spark = 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字段的DataFrame
    df = df.select(columns_to_check)

    # 创建一个临时视图
    df.createOrReplaceTempView("my_table")

    spark.sql("select DISTINCT * from my_table").show()

    # 停止SparkSession
    spark.stop()

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

import os
from pyspark.sql import SparkSession

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'
    os.environ['PYSPARK_DRIVER_PYTHON'] = 'E:/penv/pyspark_env/python.exe'

    # 创建SparkSession
    spark = 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 tt
        WHERE ID  IN (select min(ID) from tt group by Name,Email,Phone,Address)
    """).show()

    # 停止SparkSession
    spark.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编程实战

相关文章:

  • Flask + Pear Admin Layui 快速开发管理后台
  • 隐式类型转换
  • Hyperlane 框架路由功能详解:静态与动态路由全掌握
  • 无锁队列简介与实现示例
  • # 基于人脸关键点的多表情实时检测系统
  • 4月7号.
  • 【开源宝藏】30天学会CSS - DAY12 第十二课 从左向右填充的文字标题动画
  • spring-cloud-alibaba-nacos-discovery使用说明
  • 超大规模数据场景(思路)——面试高频算法题目
  • 进程和线程的区别和联系
  • 【Java面试系列】Spring Boot应用中的事务传播机制与分布式事务实践详解 - 3-5年Java开发必备知识
  • 【软件】在 macOS 上安装和配置 Apache HTTP 服务器
  • React-narice安卓打包流程
  • ifconfig 使用详解
  • animals_classification动物分类
  • 子类是否能继承
  • 解决windows下删除文件提示该项目不存在
  • 设计模式简述(七)原型模式
  • Qt音频采集:QAudioInput详解与示例
  • Android打包及上架应用市场问题处理
  • 优设/seo网站优化软件
  • 零基础学习网站建设/搜索引擎优化的含义
  • 国土分局网站建设方案/中国十大流量网站
  • 如何做网店网站/张文宏说上海可能是疫情爆发
  • 网站建设 英语词汇/百度一下就一个
  • 本溪市做网站公司/沈阳头条今日头条新闻最新消息