【发布时间】:2015-09-08 02:18:33
【问题描述】:
我在 Hadoop 的 YARN 上运行 Spark。这种转换是如何工作的?转换前是否发生了 collect()?
我还需要在每个从节点上安装 Python 和 R 才能进行转换吗?我正在努力寻找这方面的文档。
【问题讨论】:
标签: pandas apache-spark dataframe hadoop apache-spark-sql
我在 Hadoop 的 YARN 上运行 Spark。这种转换是如何工作的?转换前是否发生了 collect()?
我还需要在每个从节点上安装 Python 和 R 才能进行转换吗?我正在努力寻找这方面的文档。
【问题讨论】:
标签: pandas apache-spark dataframe hadoop apache-spark-sql
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.DataFrame 和base::data.frame 分别在 Python 和 R 中)。
矢量化用户定义函数
由于Spark 2.3.0 PySpark 还提供了一组pandas_udf(SCALAR、GROUPED_MAP、GROUPED_AGG),它们在由
SCALAR 变体的分区GROUPED_MAP 和 GROUPED_AGG 的分组表达式。每个块由
表示pandas.core.series.Series,如果是SCALAR 和GROUPED_AGG 变体。pandas.core.frame.DataFrame 的单个 GROUPED_MAP 变体。同样,从 Spark 2.0.0 开始,SparkR 提供了dapply 和gapply 函数,分别在由分区和分组表达式定义的data.frames 上运行。
上述功能:
coalesce(1))或分组表达式很简单(即groupBy(lit(1))),否则不存在单节点瓶颈。【讨论】: