大规模 PySpark

PySpark 入门

Benjamin Schmidt

Data Engineer

利用规模优势

  • PySpark 可高效处理 GB 至 TB 级数据
  • 目标是速度与高效处理
  • 理解执行机制可进一步提效
  • 使用 broadcast 协调整个集群
joined_df = large_df.join(broadcast(small_df), 
                          on="key_column", how="inner")
joined_df.show()
PySpark 入门

执行计划

# Using explain() to view the execution plan
df.filter(df.Age > 40).select("Name").explain()
== Physical Plan ==
*(1) Filter (isnotnull(Age) AND (Age > 30))
+- Scan ExistingRDD[Name:String, Age:Int]
1 https://spark.apache.org/docs/latest/api/python/reference/pyspark.sql/api/pyspark.sql.DataFrame.explain.html
PySpark 入门

缓存与持久化 DataFrame

  • 缓存:将数据存入内存,适合小数据集,访问更快
  • 持久化:按不同存储级别保存,适合更大数据集
df = spark.read.csv("large_dataset.csv", header=True, inferSchema=True)

# Cache the DataFrame
df.cache()

# Perform multiple operations on the cached DataFrame df.filter(df["column1"] > 50).show() df.groupBy("column2").count().show()
PySpark 入门

按不同存储级别持久化 DataFrame

# Persist the DataFrame with storage level
from pyspark import StorageLevel

df.persist(StorageLevel.MEMORY_AND_DISK)

# Perform transformations result = df.groupBy("column3").agg({"column4": "sum"}) result.show() # Unpersist after use df.unpersist()
PySpark 入门

优化 PySpark

  • 小范围处理:数据越多越慢;优先用 map() 等选择性更高的方法,而非 groupby()
  • 广播连接:Broadcast 会占用全部计算资源,即使在小数据集上
  • 避免重复执行:对同一数据反复执行会耗时耗算力且无收益
PySpark 入门

Passons à la pratique !

PySpark 入门

Preparing Video For Download...