【问题标题】:pyspark operations on datasets对数据集的 pyspark 操作
【发布时间】:2021-03-15 04:45:18
【问题描述】:

我最近开始学习 pySpark 进行网络图操作。我仍在尝试了解 spark 是如何使用数据帧工作的。

据我了解:

  1. 数据帧中的数据分布在一组机器上(与使用 hadoop 分布的方式相同)。

  2. 对数据帧的操作以相同的方式分布(或者,在我的情况下,本地使用给定数量的线程)。

我的主要问题是:在 pyspark 的 shell 中执行后,是仅对数据帧的操作分发还是所有操作(比如说在 python 的字典上)分发?

感谢任何进一步的澄清或文章。

【问题讨论】:

  • 通常它只是对数据帧的操作。如果您还不能使用内置聚合函数执行此操作,您可以编写 UDF 和 UDAF 以分布式方式执行自定义逻辑。
  • @Hitobat 你好!所以,如果我理解正确的话,假设我有一个包含 1000 万个键/值对的字典。如果我实现一个将字典作为参数的 UDF,则操作将以分布式方式完成,尽管字典不是数据框?
  • 您只能在数据帧上运行 udf,因此您首先需要转换输入数据。一种方法是您可以将其保存为 csv/json/avro 并使用 spark 读取它。或者直接在你的程序中使用spark.createDataFrame()或旧的spark.parallelize().toDF()进行转换。

标签: python apache-spark pyspark


【解决方案1】:

Pyspark 是用于 spark 的 python API,因此在 spark 和 python 之间有一座桥梁。它被称为Py4j,一个确保从python语言到JVM的数据序列化和编码的库,为什么?

Spark 是用 scala 编写的,这是一种类似于 java 的 POO 和 FP 语言,但不同的是,作为 Java,scala 是在 JVM(Java 虚拟机)之上工作的。与使用 scala 或 Java 相比,通过 Py4j 进行的这个过程使得通过 spark 进行数据处理的速度变慢。当谈到火花操作时,我们谈到两种类型的操作: 动作和转换:转换是一种惰性求值操作,惰性求值是什么意思?它是 spark 中的一个属性,基本上是一个 scala 特性,当表达式/值在被调用之前不被评估并且它仍然在内存中,因为这个 spark 是内存计算引擎。操作是在调用后直接评估的操作,例如 show、write、repartition....

我也需要澄清一点,spark提交模式有两种:客户端模式和集群模式。使用 spark-shell 是一种客户端模式。我邀请您详细了解这两种模式以及它们之间的区别。提交后,spark 中的任何操作都分布在纱线集群或 centos 中的所有执行器上……正如您所提到的。任何转换都不会直接评估,在动作调用 spark 创建名为 DAG 的东西后,这里是 spark DAG 的示例:(有向无环图),用于 RDD 的可视化表示和操作对他们进行。 RDD 由顶点表示,而操作由边表示。每条边都从“早期状态”指向“晚期状态”。对于每个执行者,都会分配一个任务,完成该任务后,它会通知管理器 (yarn /centos/...) 其状态等等。
您可以阅读有关 Dataframe 、 Dataset 和 RDD 的更多信息,但我可以说: RDD:弹性分布式数据集,只读模式下非结构化数据的集合(不可变 == 我们无法更改),数据帧是结构化的集合具有模式和数据集的不可变数据就像数据框,但是它是类型化的。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-01-07
    • 1970-01-01
    • 1970-01-01
    • 2017-02-11
    • 1970-01-01
    相关资源
    最近更新 更多