PySpark 中的弹性分布式数据集(RDD)

PySpark 入门

Benjamin Schmidt

Data Engineer

PySpark 中的并行化是什么?

  • 在集群的多个节点上自动并行化数据与计算
  • 在多个节点上对海量数据进行分布式处理
  • Worker 节点并行处理数据,任务结束时再汇总
  • 大规模更快(可达 GB 甚至 TB 级)

并行化

PySpark 入门

理解 RDD

RDD(Resilient Distributed Dataset,弹性分布式数据集):

  • 跨集群分布的数据集合,节点故障可自动恢复
  • 适合大规模数据
  • 不可变,可用 map()filter() 等转换;用 collect() 等动作取结果,或用 paralelize() 创建 RDD
PySpark 入门

创建 RDD

# Initialize a Spark session
from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("RDDExample").getOrCreate()

# Create a DataFrame from a csv census_df = spark.read.csv("/census.csv")
# Convert DataFrame to RDD census_rdd = census_df.rdd
# Show the RDD's contents using collect() census_rdd.collect()
PySpark 入门

演示 collect

# Collect the entire DataFrame into a local Python list of Row objects
data_collected = df.collect()

# Print the collected data
for row in data_collected:
    print(row)
```    
PySpark 入门

RDD vs DataFrame

DataFrame

  • 高层:易用且有优化
  • 类 SQL 操作:可用类 SQL 查询,少量代码完成复杂操作
  • 模式信息:含列与类型,类似 SQL 表

RDD

  • 底层:更灵活,但复杂操作需更多代码
  • 类型安全:保留数据类型,但缺乏 DataFrame 的优化
  • 无模式:处理结构化/关系型数据更困难
  • 擅长大规模扩展
  • 与 DataFrame 相比非常啰嗦,分析能力较弱
PySpark 入门

常用函数与方法

  • map(): 将函数(含自定义函数,如 lambda)应用到数据集: rdd.map(map_function)
  • collect(): 从集群收集数据: rdd.collect()
PySpark 入门

Passons à la pratique !

PySpark 入门

Preparing Video For Download...