You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过指定查询语句从MongoDB读取数据到Spark

在Spark中实现MongoDB指定查询的方案

你可以通过两种方式实现和原MongoDB查询等价的逻辑,推荐优先使用第一种(MongoDB端执行过滤投影,性能更优):

方式一:读取时直接指定过滤与投影(推荐)

通过ReadConfig传入过滤条件和投影规则,将查询逻辑推送到MongoDB端执行,只返回符合要求的数据:

import com.mongodb.spark.config.ReadConfig

val readConf = ReadConfig(Map(
  "uri" -> host,
  "database" -> "nodebb",
  "collection" -> "objects",
  // 对应MongoDB的查询过滤条件
  "filter" -> """{ "_key": { "$in": ["user:130"] } }""",
  // 对应MongoDB的投影规则:排除_id,只保留uid和username
  "projection" -> """{ "_id": 0, "uid": 1, "username": 1 }"""
))

val data = spark.read.mongo(readConf)
// 查看结果
data.show()

方式二:全量读取后在Spark端处理

如果需要后续对全量数据做更多操作,可以先读取全量数据,再用DataFrame API过滤和选择列:

import com.mongodb.spark.config.ReadConfig

val readConf = ReadConfig(Map("uri" -> host, "database" -> "nodebb", "collection" -> "objects"))
val fullData = spark.read.mongo(readConf)

// 过滤_key在指定列表中的数据,选择需要的列
val targetData = fullData
  .filter($"_key".isin("user:130"))
  .select("uid", "username")

// 查看结果
targetData.show()

注意:方式一的性能远优于方式二,尤其是当集合数据量较大时,因为它避免了全量数据的传输和加载。

内容的提问来源于stack exchange,提问作者sherin

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.13 07:15:38