如何在Databricks中从Cosmos MongoDB导入数据时添加过滤条件
解决Databricks从Cosmos MongoDB导入数据时的过滤条件错误
错误原因
你传入pipeline参数的是单一过滤文档{"type":"student"},但MongoDB聚合管道要求传入包含阶段操作的数组。这个文档被当成了一个管道阶段名称,因此触发"Unrecognized pipeline stage name: type"错误。
正确实现方式
方法1:使用聚合管道(pipeline选项)
构造包含$match阶段的聚合管道数组,这是Mongo Spark Connector推荐的过滤方式:
import json # 构造聚合管道:用$match阶段过滤type=student的数据 pipeline = [{"$match": {"type": "student"}}] df = spark.read \ .format('com.mongodb.spark.sql.DefaultSource') \ .option('uri', sourceCosmosConnectionString) \ .option('database', sourceCosmosDocument) \ .option('collection', sourceCosmosCollection) \ .option('pipeline', json.dumps(pipeline)) \ .load()
方法2:使用filter选项(更简洁)
如果你的Mongo Spark Connector版本支持,可以直接用filter选项传入过滤条件,无需构造聚合管道:
import json filter_condition = {"type": "student"} df = spark.read \ .format('com.mongodb.spark.sql.DefaultSource') \ .option('uri', sourceCosmosConnectionString) \ .option('database', sourceCosmosDocument) \ .option('collection', sourceCosmosCollection) \ .option('filter', json.dumps(filter_condition)) \ .load()
说明
两种方式都会在数据源端(Cosmos MongoDB)完成过滤,避免全量导入后再过滤,能有效提升数据导入性能。
内容的提问来源于stack exchange,提问作者Swapnil
相关产品推荐
相关产品推荐

