【问题标题】:Requirements for converting Spark dataframe to Pandas/R dataframe将 Spark 数据帧转换为 Pandas/R 数据帧的要求
【发布时间】:2015-09-08 02:18:33
【问题描述】:

我在 Hadoop 的 YARN 上运行 Spark。这种转换是如何工作的?转换前是否发生了 collect()?

我还需要在每个从节点上安装 Python 和 R 才能进行转换吗?我正在努力寻找这方面的文档。

【问题讨论】:

    标签: pandas apache-spark dataframe hadoop apache-spark-sql


    【解决方案1】:

    toPandas (PySpark) / as.data.frame (SparkR)

    必须在创建本地数据框之前收集数据。例如toPandas 方法如下所示:

    def toPandas(self):
        import pandas as pd
        return pd.DataFrame.from_records(self.collect(), columns=self.columns)
    

    您需要 Python,最好在每个节点上安装所有依赖项。

    SparkR 对应项 (as.data.frame) 只是 collect 的别名。

    总结这两种情况下的数据是collected 到驱动程序节点并转换为本地数据结构(pandas.DataFramebase::data.frame 分别在 Python 和 R 中)。

    矢量化用户定义函数

    由于Spark 2.3.0 PySpark 还提供了一组pandas_udfSCALARGROUPED_MAPGROUPED_AGG),它们在由

    定义的数据块上并行操作>
    • SCALAR 变体的分区
    • GROUPED_MAPGROUPED_AGG 的分组表达式。

    每个块由

    表示
    • 一个或多个pandas.core.series.Series,如果是SCALARGROUPED_AGG 变体。
    • pandas.core.frame.DataFrame 的单个 GROUPED_MAP 变体。

    同样,从 Spark 2.0.0 开始,SparkR 提供了dapplygapply 函数,分别在由分区和分组表达式定义的data.frames 上运行。

    上述功能:

    • 不要向司机收取。除非数据仅包含单个分区(即coalesce(1))或分组表达式很简单(即groupBy(lit(1))),否则不存在单节点瓶颈。
    • 在相应执行器的内存中加载相应的块。因此,它受到每个执行程序上可用的单个块/内存大小的限制。

    【讨论】:

    • 那么,toPandas 总是在驱动节点上?而且你永远不能在工作节点的地图函数中使用熊猫数据框?
    • @Matthias toPandas 始终在驱动程序上。如果需要,您可以在 map 中使用 pandas 对象,但这与 Spark 无关。无论你在 executor 线程中得到什么,都只是一个普通的本地对象。
    • 啊,谢谢你的澄清。我今天刚刚和同事讨论了一个场景 (see this link)。
    • 嗯,根据您的示例,它确实以分布式方式工作吗?这意味着您可以在 map 函数中创建许多 pandas 数据框。然后它们是否位于驱动程序节点上?
    猜你喜欢
    • 1970-01-01
    • 2016-09-27
    • 2020-07-24
    • 2019-01-16
    • 1970-01-01
    • 2016-01-03
    • 1970-01-01
    • 2014-04-15
    • 2021-04-19
    相关资源
    最近更新 更多