如何使用Spark Mongo连接器在Mongo查询管道中执行聚合及报错排查
问题原因
- 聚合管道语法存在错误:
$project阶段的字典闭合后多写了一个多余的右大括号,导致整个pipeline列表结构非法。 - 参数格式不符合要求:MongoDB Spark连接器的
pipeline参数仅接收JSON格式字符串,直接传入Python原生字典列表会被识别为非法聚合结构,触发报错。
正确实现方案
首先导入json模块,修正pipeline的语法错误后,将Python字典列表转换为JSON字符串再传入参数,示例代码如下:
import json # 修正语法错误,移除多余的右大括号 pipeline = [ { '$match': { 'createdDateTime': { '$gte': {'$date': f'{yesterday}T00:00:00Z'}, '$lte': {'$date': f'{today}T00:00:00Z'} } } }, { '$project': { '_class': {'$ifNull': ['$_class', '']} } } ] # 转换为JSON字符串后传入 df = spark.read.format("com.mongodb.spark.sql.DefaultSource")\ .option("spark.mongodb.input.uri", uri)\ .option("pipeline", json.dumps(pipeline))\ .load()
补充说明
如果使用3.0及以上版本的MongoDB Spark连接器,可以直接调用spark.read.mongodb()方法简化写法,无需手动指定数据源全类名。如果日期过滤出现时区偏移问题,可直接将日期字段值替换为和MongoDB存储格式对齐的ISO字符串即可。
内容的提问来源于stack exchange,提问作者Nimmala
相关产品推荐
相关产品推荐

