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

Python操作Spark常用命令指南

 更新時間:2026年02月04日 14:18:48   作者:AZ-直到世界的盡頭  
Python操作Spark的常用命令指南,涵蓋從環(huán)境配置到數(shù)據(jù)分析的核心操作,本文介紹了如何使用PySpark進(jìn)行環(huán)境配置和核心操作,包括數(shù)據(jù)讀寫、處理與轉(zhuǎn)換、聚合與高級分析,以及運(yùn)行SQL查詢和性能優(yōu)化,感興趣的朋友跟隨小編一起看看吧

Python操作Spark的常用命令指南,涵蓋從環(huán)境配置到數(shù)據(jù)分析的核心操作。

一、環(huán)境配置與Spark初始化

PySpark是Apache Spark的Python API,它結(jié)合了Python的易用性和Spark的分布式計算能力。

1. 安裝與基礎(chǔ)導(dǎo)入

# 安裝PySpark
!pip install pyspark
# 導(dǎo)入必要的庫
from pyspark.sql import SparkSession
import pyspark.sql.functions as F  # 常用函數(shù)
import pyspark.sql.types as T      # 數(shù)據(jù)類型

2. 創(chuàng)建SparkSession
SparkSession是與Spark集群交互的統(tǒng)一入口點,絕大多數(shù)操作都從這里開始。

spark = SparkSession.builder \
    .appName("MySparkApp") \    # 設(shè)置應(yīng)用名稱
    .config("spark.driver.memory", "2g") \  # 配置參數(shù)(可選)
    .getOrCreate()

3. 基礎(chǔ)環(huán)境檢查

# 檢查Spark版本
print(spark.version)
# 查看當(dāng)前配置
print(spark.conf.getAll())
# 列出數(shù)據(jù)庫和表
spark.sql("SHOW DATABASES").show()
spark.sql("USE your_database")  # 切換數(shù)據(jù)庫
spark.sql("SHOW TABLES").show()

二、核心操作:數(shù)據(jù)讀寫

1. 創(chuàng)建DataFrame
有多種方式創(chuàng)建DataFrame,這是Spark中的核心數(shù)據(jù)結(jié)構(gòu)。

# 方式1:從Python列表創(chuàng)建
data = [("Alice", 29), ("Bob", 35)]
columns = ["name", "age"]
df = spark.createDataFrame(data, schema=columns)
# 方式2:指定詳細(xì)模式(Schema)
schema = T.StructType([
    T.StructField("name", T.StringType(), True),
    T.StructField("age", T.IntegerType(), True)
])
df = spark.createDataFrame(data, schema=schema)

2. 從文件讀取數(shù)據(jù)
PySpark支持CSV、Parquet、JSON等多種格式。

# 讀取CSV文件(常用)
df = spark.read.csv(
    "path/to/file.csv",
    header=True,        # 第一行作為列名
    inferSchema=True    # 自動推斷列類型
)
# 讀取Parquet文件(列式存儲,高效)
df = spark.read.parquet("path/to/file.parquet")
# 讀取JSON文件
df = spark.read.json("path/to/file.json")

3. 數(shù)據(jù)寫入與保存

# 保存為Parquet格式(推薦,壓縮率高)
df.write.mode("overwrite").parquet("output_path.parquet")
# 保存為CSV格式
df.write.mode("overwrite") \
    .option("header", True) \
    .csv("output_path.csv")
# 保存為Spark表(可在集群中持久化)
df.write.saveAsTable("table_name")

三、數(shù)據(jù)處理與轉(zhuǎn)換

數(shù)據(jù)處理的核心是對DataFrame進(jìn)行列和行的操作。

1. 列操作

# 選擇特定列
df.select("name", "age").show()
# 創(chuàng)建新列(示例:年齡加1)
df = df.withColumn("age_plus_one", F.col("age") + 1)
# 重命名列
df = df.withColumnRenamed("old_name", "new_name")
# 更改列類型(示例:整型轉(zhuǎn)字符串)
df = df.withColumn("age_str", F.col("age").cast(T.StringType()))
# 刪除列
df = df.drop("column_to_remove")

2. 行操作(過濾與排序)

# 過濾行(示例:年齡大于30)
df_filtered = df.filter(F.col("age") > 30)
# 等價寫法
df_filtered = df.where(df["age"] > 30)
# 排序
df_sorted = df.orderBy(F.col("age").desc())  # 按年齡降序

3. 處理缺失值

# 刪除包含任何空值的行
df_clean = df.dropna()
# 填充缺失值(示例:用0填充特定列)
df_filled = df.fillna({"age": 0, "name": "Unknown"})

四、數(shù)據(jù)聚合與高級分析

1. 分組與聚合
這是數(shù)據(jù)分析中最常用的操作之一。

# 基礎(chǔ)分組聚合(示例:按部門計算平均工資)
df.groupBy("department").agg(
    F.avg("salary").alias("avg_salary"),
    F.count("*").alias("employee_count")
).show()

2. 連接(Join)操作
用于合并兩個DataFrame。

# 假設(shè)有另一個部門信息表df_dept
df_joined = df.join(
    df_dept, 
    on="department_id",  # 連接鍵
    how="inner"          # 連接方式:inner, left, right, outer等
)

3. 窗口函數(shù)
用于計算排名、移動平均等高級分析。

from pyspark.sql.window import Window
# 定義窗口:按部門分區(qū),按工資降序
window_spec = Window.partitionBy("department").orderBy(F.col("salary").desc())
# 計算部門內(nèi)工資排名
df.withColumn("salary_rank", F.rank().over(window_spec)).show()

4. 用戶自定義函數(shù)(UDF)
當(dāng)內(nèi)置函數(shù)無法滿足需求時使用。

from pyspark.sql.functions import udf
# 定義Python函數(shù)
def categorize_age(age):
    return "Young" if age < 30 else "Senior"
# 注冊為UDF(需指定返回類型)
categorize_udf = udf(categorize_age, T.StringType())
# 應(yīng)用UDF
df.withColumn("age_group", categorize_udf(F.col("age"))).show()

五、在PySpark中運(yùn)行SQL查詢

你可以直接在PySpark中執(zhí)行SQL語句,這為熟悉SQL的用戶提供了便利。

# 將DataFrame注冊為臨時視圖
df.createOrReplaceTempView("people")
# 執(zhí)行SQL查詢
result = spark.sql("""
    SELECT department, AVG(salary) as avg_sal
    FROM people 
    WHERE age > 25
    GROUP BY department
    ORDER BY avg_sal DESC
""")
result.show()

六、性能優(yōu)化與最佳實踐

  • 避免數(shù)據(jù)混洗groupByjoin等操作可能導(dǎo)致數(shù)據(jù)在節(jié)點間大量移動,應(yīng)盡量減少這類操作或提前過濾數(shù)據(jù)。
  • 選擇合適的數(shù)據(jù)格式:生產(chǎn)環(huán)境中,Parquet通常是比CSV更好的選擇,因為它支持列式存儲和謂詞下推,能顯著提高查詢性能。
  • 利用緩存:對需要多次使用的中間結(jié)果進(jìn)行緩存。
df.cache()  # 將DataFrame緩存到內(nèi)存
df.unpersist()  # 使用后釋放緩存
  • 及時關(guān)閉會話:處理完成后關(guān)閉SparkSession以釋放資源。
spark.stop()

七、一個完整的示例

from pyspark.sql import SparkSession
import pyspark.sql.functions as F
# 1. 初始化
spark = SparkSession.builder.appName("Example").getOrCreate()
# 2. 讀取數(shù)據(jù)
df = spark.read.csv("sales.csv", header=True, inferSchema=True)
# 3. 數(shù)據(jù)處理
result = (df
    .filter(F.col("amount") > 100)           # 篩選大額交易
    .groupBy("region", "product")            # 按地區(qū)和產(chǎn)品分組
    .agg(F.sum("amount").alias("total_sales")) # 計算總銷售額
    .orderBy(F.col("total_sales").desc())    # 按銷售額降序排序
)
# 4. 輸出
result.show()
result.write.mode("overwrite").parquet("sales_summary.parquet")
# 5. 清理
spark.stop()

到此這篇關(guān)于Python操作Spark常用命令指南的文章就介紹到這了,更多相關(guān)python spark命令內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Python根據(jù)字符串語義相似度計算的多種實現(xiàn)算法

    Python根據(jù)字符串語義相似度計算的多種實現(xiàn)算法

    這篇文章主要介紹了Python根據(jù)字符串語義相似度計算的多種實現(xiàn)算法,以下是幾種基于語義的字符串相似度計算方法,每種方法都會返回0.0到1.0之間的相似度分?jǐn)?shù)(保留一位小數(shù)),需要的朋友可以參考下
    2025-07-07
  • 一文帶你搞懂Python如何解析CSV文件

    一文帶你搞懂Python如何解析CSV文件

    CSV(Comma-Separated?Values)是一種以純文本格式存儲表格數(shù)據(jù)的文件格式,本文主要為大家詳細(xì)介紹了如何使用Python解析CSV文件,有需要的小伙伴可以了解下
    2026-03-03
  • python提取圖像的名字*.jpg到txt文本的方法

    python提取圖像的名字*.jpg到txt文本的方法

    下面小編就為大家分享一篇python提取圖像的名字*.jpg到txt文本的方法,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2018-05-05
  • Django實現(xiàn)表單驗證

    Django實現(xiàn)表單驗證

    這篇文章主要為大家詳細(xì)介紹了Django實現(xiàn)表單驗證的相關(guān)資料,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2018-09-09
  • Python3.5編程實現(xiàn)修改IIS WEB.CONFIG的方法示例

    Python3.5編程實現(xiàn)修改IIS WEB.CONFIG的方法示例

    這篇文章主要介紹了Python3.5編程實現(xiàn)修改IIS WEB.CONFIG的方法,涉及Python針對xml格式文件的讀寫以及節(jié)點操作相關(guān)技巧,需要的朋友可以參考下
    2017-08-08
  • 詳解Python遍歷列表時刪除元素的正確做法

    詳解Python遍歷列表時刪除元素的正確做法

    這篇文章主要介紹了詳解Python遍歷列表時刪除元素的正確做法,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2021-01-01
  • Python萬字深入內(nèi)存管理講解

    Python萬字深入內(nèi)存管理講解

    內(nèi)存管理是指在程序的運(yùn)行過程中,分配內(nèi)容和回收內(nèi)存的過程。如果只分配,不回收,電腦上那點內(nèi)存很快就被用光。幸運(yùn)的是,Python和Java等高級語言會自動管理內(nèi)存的分配和回收
    2022-07-07
  • Python實現(xiàn)switch/case語句

    Python實現(xiàn)switch/case語句

    與Java、C\C++等語言不同,Python中是不提供switch/case語句的,這一點讓我感覺到很奇怪。我們可以通過如下幾種方法來實現(xiàn)switch/case語句
    2021-08-08
  • python自動發(fā)郵件總結(jié)及實例說明【推薦】

    python自動發(fā)郵件總結(jié)及實例說明【推薦】

    python發(fā)郵件需要掌握兩個模塊的用法,smtplib和email,這倆模塊是python自帶的,只需import即可使用。這篇文章主要介紹了python自動發(fā)郵件總結(jié)及實例說明 ,需要的朋友可以參考下
    2019-05-05
  • pip下載安裝時使用國內(nèi)源配置方式

    pip下載安裝時使用國內(nèi)源配置方式

    這篇文章主要介紹了pip下載安裝時使用國內(nèi)源配置方式,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2026-04-04

最新評論

黄梅县| 潮州市| 富民县| 毕节市| 津市市| 临桂县| 阳谷县| 边坝县| 莱西市| 舒城县| 孝义市| 神池县| 平谷区| 邢台县| 凌源市| 杭锦旗| 阿瓦提县| 普安县| 丹凤县| 商丘市| 浦东新区| 蒲城县| 泰顺县| 哈密市| 韶山市| 武川县| 乐安县| 江口县| 林口县| 海晏县| 松溪县| 裕民县| 屏南县| 遵义市| 侯马市| 托里县| 双辽市| 嘉定区| 武穴市| 荣昌县| 灵山县|