如何通过指定查询语句从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
相关产品推荐
相关产品推荐

