【发布时间】:2018-04-29 03:49:44
【问题描述】:
我有一个如下所示的 MongoDB 集合:
{
"_id" : { "customerName" : "Bob", "customerPhone" : "123-456-7890"},
"purchases": ["A", "B", "C", "D"]
}
基本上,_id 是一对关于客户的唯一键,而购买是客户购买的商品的数组。
我还有一个 PySpark DataFrame,我想将它推送到这个集合中,其中包含我想更新这个特定文档的信息。
df.write.format("com.mongodb.spark.sql.DefaultSource").mode("append") \
.option("spark.mongodb.output.uri", "mongodb://localhost:27017/customer.purchases").save()
问题是,如果我要更新此文档,我想为 Bob 添加新购买,它只会在 purchases 中附加不存在的内容,而不是全部附加。
因此,我现在最终要做的是,我只需要调用 rdd.collect() 将整个内容转换为列表,而不是使用架构将其转换为 DataFrame。然后在检查密钥是否存在的同时将所有内容一一插入;当 RDD 的查询变大时,这会导致这部分变慢并且需要大量内存。
对于版本:
PySpark:2.2 MongoDB:3.0.15 Mongo Spark 连接器:2.2.1
如果我可以使用数据框将数组中的所有元素附加到 MongoDB 集合中,是否有人可以做些什么? 另外,如果我有什么遗漏或其他我应该做的事情,请告诉我。 谢谢!
【问题讨论】:
标签: apache-spark pyspark spark-dataframe pymongo pyspark-sql