本文作为Data Science Blogathon的一部分发表。
介绍
在本教程中,我们将学习PySpark的构建块,称为弹性分布式数据集,也称为PySparkRDD。在我们这样做之前,让我们先了解一下它的基本概念。

什么是RDD?
RDD代表弹性分布式数据集,它是在多个节点上运行和工作以在集群中执行并行处理的元素。RDD是不可变的,这意味着一旦创建了RDD,就不能更改它们。它们是容错的,因此在发生故障时会自动恢复。您可以对这些RDD应用多个操作以实现特定任务。
有两种方法可以应用操作:
• 转型
• 事件
让我们看看这些方法是什么:
转型— 这些操作用于创建新的RDD。过滤器、组和映射是转换的示例。
事件— 这些操作应用于RDD,指示Spark执行计算并将结果发送回控制器。
要在PySpark中使用任何操作,我们需要首先创建PySparkRDD。下面的代码块详细介绍了PySpark RDD—
class pyspark.RDD (
Judd,
ctx
jrdd_deserializer = AutoBatchedSerializer(PickleSerializer())
)
让我们看看如何使用PySpark运行一些基本操作。Python文件中的以下代码创建了一个单词RDD,用于存储一组所述单词。
words = sc.parallelize (
["scale",
"Java",
"Hadoop",
"spark",
"Akka",
"spark vs Hadoop",
"pyspark",
"pyspark and spark"]
)
现在让我们做一些单词运算。
计数()
将返回RDD中的元素数组。
---------------------------------------count.py------- --------------------------------
from pyspark import SparkContext
sc = SparkContext("local", "count app")
words = sc.parallelize (
["scale",
"Java",
"Hadoop",
"spark",
"Akka",
"spark vs hadoop",
"pyspark",
"pyspark and spark"]
)
counts = words.count()
print "elements are -> %i" % (counts)
---------------------------------------count.py------- --------------------------------
命令− count()的命令是—
$SPARK_HOME/bin/spark-submit count.py
输出− 上述代码的输出—
Number of elements in RDD → 8
收集()
返回所有元素。
---------------------------------------collect.py------- --------------------------------
from pyspark import SparkContext
sc = SparkContext("local", "Collect app")
words = sc.parallelize (
["scale",
"Java",
"Hadoop",
"spark",
"Akka",
"spark vs hadoop",
"pyspark",
"pyspark and spark"]
)
coll = words.collect()
print "RDD elements -> %s" % (collection)
---------------------------------------collect.py------- --------------------------------
命令− collect()的命令是—
$SPARK_HOME/bin/spark-submit collect.py
输出− 上述代码的输出—
Elements in RDD -> [
'scale',
'Java',
'hadoop',
'spark',
'Akka',
'spark vs Hadoop,
'pyspark',
'pyspark and spark'
]
foreach(f)
仅返回满足foreach内函数条件的元素。在下面的示例中,我们调用print In foreach函数,该函数打印RDD中的所有元素。
----------------------------------------foreach.py------- --------------------------------
from pyspark import SparkContext
sc = SparkContext("local", "ForEach app")
words = sc.parallelize (
["scale",
"Java",
"Hadoop",
"spark",
"Akka",
"spark vs hadoop",
"pyspark",
"pyspark and spark"]
)
def f(x): print(x)
fore = words.foreach(f)
----------------------------------------foreach.py------- --------------------------------
命令− 这是foreach(f)的命令—
$SPARK_HOME/bin/spark-submit foreach.py
输出− 上述代码的输出—
scala
Java
hadoop
spark
Akka
spark vs hadoop
pyspark
pyspark and spark
过滤器(f)
返回一个新的RDD,其中包含满足过滤器内函数的元素。在下面的示例中,我们过滤掉包含“spark”的字符串
----------------------------------------filter.py------- --------------------------------
from pyspark import SparkContext
sc = SparkContext("local", "Filter app")
words = sc.parallelize (
["scale",
"Java",
"Hadoop",
"spark",
"Akka",
"spark vs hadoop",
"pyspark",
"pyspark and spark"]
)
word_filter = words.filter(lambda x: 'spark' in x)
filtered = word_filter.collect()
print "Filtered RDD -> %s" % (filtered)
----------------------------------------filter.py------- ---------------------------------
命令− 过滤器(f)命令为—
$SPARK_HOME/bin/spark-submit filter.py
输出− 上述代码的输出—
Filtered RDD -> [
'spark',
'spark vs Hadoop,
'pyspark',
'pyspark and spark'
]
映射(f,保留分区=False)
通过将函数应用于其中的每个元素,将返回一个新的RDD。在下面的示例中,我们创建一个键值对,并将每个字符串映射到值1。
----------------------------------------map.py------- --------------------------------
from pyspark import SparkContext
sc = SparkContext("local", "Map app")
words = sc.parallelize (
["scale",
"Java",
"Hadoop",
"spark",
"Akka",
"spark vs hadoop",
"pyspark",
"pyspark and spark"]
)
word_map = words.map(lambda x: (x, 1))
mapping = word_map.collect()
print "Key-value pair -> %s" % (mapping)
----------------------------------------map.py------- --------------------------------
命令− map(f,protectspationing=False)的命令是—
$SPARK_HOME/bin/spark-submit map.py
输出− 上述代码的输出—
Key-value pair -> [
('scale', 1),
('java', 1),
('Hadoop, 1),
('spark', 1),
('akka', 1),
('spark vs hHadoop 1),
('pyspark', 1),
('pyspark and spark', 1)
]
减少(f)
执行指定的关联和交换二进制操作后,返回RDD中的元素。在下面的示例中,我们从操作符导入add包并将其应用于“num”以执行简单的加法操作。
-----------------------------------------reduce.py------- --------------------------------
from pyspark import SparkContext
from the import add operator
sc = SparkContext("local", "Reduce app")
nums = sc.parallelize([1, 2, 3, 4, 5])
add = nums.reduce(add)
print "Adding elements -> %i" % (add)
-----------------------------------------reduce.py------- --------------------------------
命令− 调用ce(f)的命令是—
$SPARK_HOME/bin/spark-submit reduction.py
输出− 上述代码的输出—
Adding all elements -> 15
加入(其他,numPartitions=None)
返回一个RDD,其中包含一对具有相应键和所有值的元素对于该特定键。下面的示例显示了两个不同RDD中的元素对。在连接这两个RDD之后,我们得到一个RDD,其中包含具有相应键及其值的元素。
----------------------------------------join.py------- --------------------------------
from pyspark import SparkContext
sc = SparkContext("local", "Join app")
x = sc.parallelize([("spark", 1), ("hadoop", 4)])
y = sc.parallelize([("spark", 2), ("Hadoop", 5)])
joined = x.join(y)
final = join.collect()
print "connect RDD -> %s" % (final)
----------------------------------------join.py------- --------------------------------
命令− 用于连接的命令(其他,numPartitions=None)是—
$SPARK_HOME/bin/spark-submit join.py
输出− 上述代码的输出—
Connect to RDD -> [
('spark', (1, 2)),
('Hadoop, (4, 5))
]
缓存()
将此RDD保持为默认存储级别(MEMORY_ONLY)。您还可以检查它是否缓存。
---------------------------------------cache.py------- --------------------------------
from pyspark import SparkContext
sc = SparkContext("local", "Cache app")
words = sc.parallelize (
["scale",
"Java",
"Hadoop",
"spark",
"Akka",
"spark vs hadoop",
"pyspark",
"pyspark and spark"]
)
words.cache()
caching = words.persist().is_cached
print "cached words are > %s" % (caching)
---------------------------------------cache.py------- --------------------------------
命令− cache()的命令是—
$SPARK_HOME/bin/spark-submit cache.py
输出− 上述代码的输出—
Words have been cached -> True
这是一些最重要的操作。
结论
RDD代表弹性分布式数据集,它是在多个节点上运行和工作以在集群中执行并行处理的元素。RDD是不可变的,这意味着一旦创建了RDD,就不能更改它们。
-
转型— 这些操作应用于RDD以创建新的RDD。过滤器、组和映射是转换的示例。
-
事件— 这些操作应用于RDD,指示Spark执行计算并将结果发送回控制器。
-
它只返回满足foreach内函数条件的元素。在下面的示例中,我们调用print In foreach函数,该函数打印RDD中的所有元素。
本文中显示的媒体并非Analytics Vidhya所有,由作者自行决定使用。
原文标题:Create RDD in Apache Spark using Pyspark
原文作者:Gitesh Dhore
原文链接:https://www.analyticsvidhya.com/blog/2022/08/create-rdd-in-apache-spark-using-pyspark/




