【发布时间】: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,但实际上我更普遍地遇到了这种情况,在不仅仅是模型数据+预测的情况下。
我的问题
- 特别是在上述情况下,如何正确设置分区以使其正常工作?
-
应该如何按行索引连接两个数据框?
- (我知道标准响应是“你不应该...分区使索引变得毫无意义”,但在 Spark 创建不会像我在上面的链接中描述的那样强制数据丢失的 ML 库之前,这将始终是一个问题.)
【问题讨论】:
-
您的数据集有多大?如果它们不是太大,一个低技术的方法是将它们都写入csv文件并使用paste将它们组合起来。
-
它们很大,但不会太大而无法放入 CSV...只要我逐行阅读它们,听起来就像 paste 一样。但这似乎与使用 Spark 的目的背道而驰。如果没有更优雅的,我一定会考虑的。
标签: python dataframe apache-spark pyspark