以下是 PySpark 中 parallelize 和 textFile 两种创建 RDD 的语法及示例:

1. parallelize:从内存集合创建 RDD

功能

将 Python 本地集合(如列表、元组等)并行化到 Spark 集群,生成分布式 RDD。

语法

python

# 基本语法
rdd = sc.parallelize(c, numSlices=None)
  • sc:SparkContext 对象(PySpark 程序的入口)。
  • c:Python 本地集合(如 [1,2,3]("a","b","c") 等可迭代对象)。
  • numSlices(可选):指定 RDD 的分区数,默认值由 Spark 配置 spark.default.parallelism 决定(通常为集群核心总数)。
示例

python

运行

from pyspark import SparkContext

# 初始化 SparkContext
sc = SparkContext("local", "ParallelizeExample")  # local 表示本地模式

# 从列表创建 RDD(默认分区数)
list_rdd = sc.parallelize([1, 2, 3, 4, 5])

# 从元组创建 RDD,指定 2 个分区
tuple_rdd = sc.parallelize(("spark", "pyspark", "hadoop"), numSlices=2)

# 查看 RDD 分区数
print(list_rdd.getNumPartitions())  # 输出:默认值(如本地模式可能为 1 或 CPU 核心数)
print(tuple_rdd.getNumPartitions())  # 输出:2

2. textFile:从文本文件创建 RDD

功能

读取文本文件(本地文件或分布式文件系统如 HDFS、S3 等),每行作为 RDD 的一个元素,生成 RDD[str]

语法

python

运行

# 基本语法
rdd = sc.textFile(name, minPartitions=None, use_unicode=True)
  • sc:SparkContext 对象。
  • name:文件路径,支持:
    • 本地文件:需用 file:// 协议(如 file:///home/user/data.txt),本地模式下可直接读取,集群模式需确保所有节点有该文件。
    • 分布式文件:如 HDFS 路径(hdfs://namenode:9000/user/data.txt)、S3 路径(s3a://bucket/data.txt)等,支持通配符(如 *.txt)或目录(读取目录下所有文件)。
  • minPartitions(可选):指定最小分区数,默认根据文件大小和块大小自动计算(如 HDFS 块默认 128MB,一个块对应一个分区)。
  • use_unicode(可选):是否将文件内容解码为 Unicode 字符串(默认 True),设为 False 则返回字节串(bytes 类型)。
示例

python

运行

from pyspark import SparkContext

sc = SparkContext("local", "TextFileExample")

# 读取本地文本文件(本地模式)
local_rdd = sc.textFile("file:///C:/data/test.txt")  # Windows 路径
# 或 Linux 本地路径:file:///home/user/test.txt

# 读取 HDFS 文本文件
hdfs_rdd = sc.textFile("hdfs://namenode:9000/user/hadoop/logs/access.log")

# 读取目录下所有 .txt 文件,指定最小分区数为 5
dir_rdd = sc.textFile("/user/data/*.txt", minPartitions=5)

# 读取非 Unicode 编码文件(返回字节串)
bytes_rdd = sc.textFile("file:///data/binary.txt", use_unicode=False)

# 查看前 5 行数据
print(local_rdd.take(5))  # 输出:文件前 5 行的字符串列表

核心说明

  • parallelize:依赖本地内存集合,数据量不宜过大(受限于 Driver 内存),常用于测试。
  • textFile:依赖外部文件系统,支持大规模分布式数据,是生产环境处理文本数据的常用方式。
  • 两者均为惰性操作:创建 RDD 时仅记录元信息,实际数据读取和计算在调用行动算子(如 takecountcollect)时触发。

更多推荐