【问题标题】:pySpark convert result of mapPartitions to spark DataFramepySpark 将 mapPartitions 的结果转换为 spark DataFrame
【发布时间】:2019-12-10 23:13:04
【问题描述】:

我有一项工作需要在分区的 spark 数据帧上运行,过程如下所示:

rdd = sp_df.repartition(n_partitions, partition_key).rdd.mapPartitions(lambda x: some_function(x))

结果是pandas.dataframe 的rdd,

type(rdd) => pyspark.rdd.PipelinedRDD
type(rdd.collect()[0]) => pandas.core.frame.DataFrame

和rdd.glom().collect() 返回结果如下:

[[df1], [df2], ...]

现在我希望将结果转换为spark dataframe,我做的方式是:

sp = None
for i, partition in enumerate(rdd.collect()):
    if i == 0:
        sp = spark.createDataFrame(partition)
    else:
        sp = sp.union(spark.createDataFrame(partition))

return sp

但是,结果可能很大,rdd.collect() 可能会超出驱动程序的内存,所以我需要避免collect() 操作。有没有办法解决这个问题?

提前致谢!

【问题讨论】:

  • 你可以运行rdd.toDF()。或者,spark.createDataFrame(rdd)
  • @samakart,不是真的,它会导致错误ValueError: The truth value of a DataFrame is ambiguous. Use a.empty, a.bool(), a.item(), a.any() or a.all().,我猜它只适用于Row。
  • 不,您必须提交正确的架构。在你的情况下,它只是无法弄清楚我想的类型。 createDataFrame 有 schema 参数。它还接受类似 sql 的字符串。我还没有在任何地方找到这种语法的文档。但它只是sql标准。
  • @dre-hh 是 this 你在找什么?
  • 是和否:)。所以是的,这是将模式提供为 python 类型的一种方法。但 schema 参数也接受更短的 sql dsl 表示法。例如。 createDataFrame(x, schema="uuid_id STRING, url STRING, title STRING")我试图在文档中找到它支持的类型,但这些基本上是 SQL 表示法中 python 数据类型的类比

标签: python apache-spark pyspark


【解决方案1】:

如果你想继续使用 rdd api。 mapPartitions 接受一个类型的迭代器并期望另一个类型的迭代器作为结果。 pandas_df 不是 mapPartitions 可以直接处理的迭代器类型。如果你必须使用 pandas api,你可以从 pandas.iterrows 创建一个合适的生成器

这样,您的整体 mapPartitions 结果将是您的行类型的单个 rdd,而不是 pandas 数据帧的 rdd。这样的 rdd 可以通过 on-the-fly 模式发现无缝转换为数据帧

from pyspark.sql import Row

def some_fuction(iter):
  pandas_df = some_pandas_result(iter)
  for index, row in pandas_df.iterrows():
     yield Row(id=index, foo=row['foo'], bar=row['bar'])


rdd = sp_df.repartition(n_partitions, partition_key).rdd.mapPartitions(lambda x: some_function(x))
df = spark.createDataFrame(rdd)

【讨论】:

  • 所以这些方法需要提前知道生成的pandas数据框的列?
  • 感谢您花时间研究这个问题。我按照您的建议解决了这个问题:转换为Row,然后转换为createDataFrame。我应用的代码是将pandas data frame 的每一行附加到Row 对象列表中:row_list.append(Row(**row_dict))
  • 很高兴它也可以像这样动态地工作。我不确定。一天也给 pandas udf(下面的答案)一个旋转。缺点是,您需要指定行类型的正确类型。但它将使用apache airflow,使用现代 cpu SIMD 指令,并且从 pandas 到数据帧的转换将在幕后以更优化的代码方式进行
【解决方案2】:

您可以直接在 datframe 上使用新的 pandas grouped udf 而不是 rdd.mapPartitions 。该函数本身接受一个组作为 pandas df 并返回 pandas df。

与spark dataframe apply api一起使用时,spark会自动将分区的pandas dataframes组合成一个新的spark 数据框。

# a grouped pandas_udf receives the whole group as a pandas dataframe
# it must also return a pandas dataframe
# the first schema string parameter must describe the return dataframe schema

# in this example the result dataframe contains 2 columns id and value
@pandas_udf("id long, value double", PandasUDFType.GROUPED_MAP)
def some_function(pdf):
    return pdf.apply(some_pdf_func)

df.groupby(df.partition_key).apply(some_function).show()

【讨论】:

    【解决方案3】:

    你可以这样做:

    sp = None 
    def f(x):
     sp = spark.createDataFrame(x)
     return (sp)
    sp = sp.union(rdd.foreach(f))
    

    参考:

    Spark SQL DataFrame

    Spark RDD

    如果可行,请投票

    【讨论】:

    • 您好,感谢您的深入了解。恐怕它不起作用。 sp=None => 'NoneType' object has no attribute 'union' ; foreach => SparkContext can only be used on the driver, not in code that it run on workers
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-03-17
    • 1970-01-01
    • 2022-06-11
    • 2022-12-18
    • 2016-05-29
    • 1970-01-01
    相关资源
    最近更新 更多