【问题标题】:Unable to write PySpark Dataframe created from two zipped dataframes无法写入从两个压缩数据帧创建的 PySpark 数据帧
【发布时间】:2020-12-22 22:05:19
【问题描述】:

我正在尝试关注the example given here 组合两个数据帧没有共享连接键(通过数据库表或熊猫数据帧中的“索引”组合,除了 PySpark 没有这个概念) :

我的代码

left_df = left_df.repartition(right_df.rdd.getNumPartitions()) # FWIW, num of partitions = 303
joined_schema = StructType(left_df.schema.fields + right_df.schema.fields)
interim_rdd = left_df.rdd.zip(right_df.rdd).map(lambda x: x[0] + x[1])
full_data = spark.createDataFrame(interim_rdd, joined_schema)

这一切似乎都很好。我在使用 DataBricks 时对其进行了测试,我可以毫无问题地运行上面的“单元格”。但是当我去保存它时,我无法,因为它抱怨分区不匹配(???)。我已经确认分区的数量匹配,但您也可以在上面看到我明确确保它们匹配。我的保存命令:

full_data.write.parquet(my_data_path, mode="overwrite")

错误

我收到以下错误:

Caused by: org.apache.spark.SparkException: Can only zip RDDs with same number of elements in each partition

我的猜测

我怀疑问题在于,即使我匹配了 number 个分区,但每个分区中的行数并不相同。但我不知道该怎么做。我只知道如何指定分区数,不知道分区方式。

或者,更具体地说,我不知道如何指定如何分区如果没有我可以使用的列。请记住,它们没有共享列。


我怎么知道我可以通过这种方式组合它们,而无需共享连接键?在这种情况下,是因为我正在尝试join model predictions with input data,但实际上我更普遍地遇到了这种情况,在不仅仅是模型数据+预测的情况下。

我的问题

  1. 特别是在上述情况下,如何正确设置分区以使其正常工作?
  2. 应该如何按行索引连接两个数据框?
    • (我知道标准响应是“你不应该...分区使索引变得毫无意义”,但在 Spark 创建不会像我在上面的链接中描述的那样强制数据丢失的 ML 库之前,这将始终是一个问题.)

【问题讨论】:

  • 您的数据集有多大?如果它们不是太大,一个低技术的方法是将它们都写入csv文件并使用paste将它们组合起来。
  • 它们很大,但不会太大而无法放入 CSV...只要我逐行阅读它们,听起来就像 paste 一样。但这似乎与使用 Spark 的目的背道而驰。如果没有更优雅的,我一定会考虑的。

标签: python dataframe apache-spark pyspark


【解决方案1】:

您可以暂时切换到 RDD 并使用zipWithIndex 添加索引。然后可以将此索引用作连接标准:

#create rdds with an additional index
#as zipWithIndex adds the index as second column, we have to switch
#the first and second column
left = left_df.rdd.zipWithIndex().map(lambda a: (a[1], a[0]))
right= right_df.rdd.zipWithIndex().map(lambda a: (a[1], a[0]))

#join both rdds 
joined = left.fullOuterJoin(right)

#restore the original columns
result = spark.createDataFrame(joined).select("_2._1.*", "_2._2.*")

zipWithIndex 的 Javadoc 声明

某些 RDD,例如 groupBy() 返回的 RDD,不保证分区中元素的顺序。

根据原始数据集的性质,此代码可能不会产生确定性结果。

【讨论】:

  • 谢谢!这是有道理的,但我担心不确定性。 :(
  • 接受我的回答,这将是确定性的。
  • 我实际上都尝试过,但由于确定性,我专注于@thebluephantom。但我无法让它工作。请参阅我在对他的评论中发布的错误。感谢大家的帮助!
  • @MikeWilliamson 确定性仍然基于行位置。我观察到 zipWI 分配到同一个分区,将 int 重新分区作为键、值对的键,并将传入的数据留在同一个位置,因为它只是一个狭窄的转换。也就是说,您必须在分区中拥有相同数量的数据,否则您没有压缩 RDD(不是 DF)的用例。
【解决方案2】:

RDD 是旧帽子,但从这个角度回答错误。

来自拉筹伯大学http://homepage.cs.latrobe.edu.au/zhe/ZhenHeSparkRDDAPIExamples.html#zip以下:

通过将任一分区的第 i 个与每个分区的第 i 个相结合来连接两个 RDD 其他。生成的 RDD 将由两个组件的元组组成 由提供的方法解释为键值对 PairRDDFunctions 扩展。

笔记对。

这意味着你必须有相同的分区器,分区数和每个分区的 kv 数,否则上面的定义不成立。

最好从文件中读取,因为 repartition(n) 可能不会给出相同的分布。

解决这个问题的一个小技巧是使用 zipWithIndex 作为 k 的 k,v,就像这样(Scala 不是 pyspark 特定的方面):

val rddA = sc.parallelize(Seq(
  ("ICCH 1", 10.0), ("ICCH 2", 10.0), ("ICCH 4", 100.0), ("ICCH 5", 100.0)
))
val rddAA = rddA.zipWithIndex().map(x => (x._2, x._1)).repartition(5)

val rddB = sc.parallelize(Seq(
  (10.0, "A"), (64.0, "B"), (39.0, "A"), (9.0, "C"), (80.0, "D"), (89.0, "D")
))
val rddBB = rddA.zipWithIndex().map(x => (x._2, x._1)).repartition(5)

val zippedRDD = (rddAA zip rddBB).map{ case ((id, x), (y, c)) => (id, x, y, c) }
zippedRDD.collect

然后重新分区(n)似乎可以工作,因为 k 是相同的类型。

但是每个分区必须有相同的 num 个元素。它就是这样,但它是有道理的。

【讨论】:

  • 我无法让它工作。总是有人抱怨分区不匹配。多次试验后的最新错误:File "/databricks/spark/python/pyspark/sql/types.py", line 1387, in verify_struct "length of fields (%d)" % (len(obj), len(verifiers)))) ValueError: Length of object (2) does not match with length of fields (31)
猜你喜欢
  • 2023-03-24
  • 2023-01-27
  • 2019-10-31
  • 2020-10-06
  • 2018-06-24
  • 1970-01-01
  • 1970-01-01
  • 2020-09-18
  • 1970-01-01
相关资源
最近更新 更多