Spark Scala代码中MongoDB Pipeline过滤选项不生效求助
问题原因与解决方案
你的问题出在MongoDB聚合管道的日期写法上:直接在字符串化的pipeline里使用ISODate()构造函数,Spark MongoDB Connector无法将其解析为MongoDB的日期类型,会被当作普通函数字面量处理,导致$match条件完全失效,所以读取的是全量数据。
修复步骤
1. 修正聚合管道的日期格式
使用MongoDB JSON规范中的$date操作符定义日期,替换原有的ISODate()写法:
val pipeline = """[ { "$match": { "updated_at": { "$gte": { "$date": "2022-03-07T03:39:19.416+00:00" } } } } ]""" val df = spark.read.format("mongodb") .option("connection.uri", connectionUri) .option("database", database) .option("collection", collection) .option("pipeline", pipeline) .load()
2. 验证字段类型
确认MongoDB集合中的updated_at字段是日期类型(Date),而非字符串类型。如果是字符串,需先将其转换为日期再比较,或调整匹配条件为字符串比较(不推荐,需保证日期字符串格式完全一致)。
备选方案:使用filter选项替代pipeline
如果只是简单过滤,也可以直接用Spark MongoDB Connector的filter选项,写法更简洁且不易出错:
import java.time.LocalDateTime import java.time.format.DateTimeFormatter // 构造过滤日期(示例为昨日) val yesterday = LocalDateTime.now().minusDays(1).format(DateTimeFormatter.ISO_DATE_TIME) val filter = s"""{"updated_at": {"$$gte": {"$$date": "$yesterday"}}}""" val df = spark.read.format("mongodb") .option("connection.uri", connectionUri) .option("database", database) .option("collection", collection) .option("filter", filter) .load()
验证
执行df.count()后,对比MongoDB中手动执行以下查询的结果数,确认数据量一致:
db.collection.find({ updated_at: { $gte: ISODate("2022-03-07T03:39:19.416+00:00") } }).count()
内容的提问来源于stack exchange,提问作者Aditya Mishra
相关产品推荐
相关产品推荐

