暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

使用Pyspark在Apache Spark中创建RDD

原创 eternity 2022-08-23
585

本文作为Data Science Blogathon的一部分发表。

介绍

在本教程中,我们将学习PySpark的构建块,称为弹性分布式数据集,也称为PySparkRDD。在我们这样做之前,让我们先了解一下它的基本概念。
image.png

什么是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/

「喜欢这篇文章,您的关注和赞赏是给作者最好的鼓励」
关注作者
【版权声明】本文为墨天轮用户原创内容,转载时必须标注文章的来源(墨天轮),文章链接,文章作者等基本信息,否则作者和墨天轮有权追究责任。如果您发现墨天轮中有涉嫌抄袭或者侵权的内容,欢迎发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论