【问题标题】:How to query to mongo using spark?如何使用 spark 查询 mongo?
【发布时间】:2015-02-15 20:50:27
【问题描述】:

我正在使用 spark 和 mongo。我可以使用以下代码连接到 mongo:

val sc = new SparkContext("local", "Hello from scala")

val config = new Configuration()
config.set("mongo.input.uri", "mongodb://127.0.0.1:27017/dbName.collectionName")
val mongoRDD = sc.newAPIHadoopRDD(config, classOf[com.mongodb.hadoop.MongoInputFormat], classOf[Object], classOf[BSONObject])

上面的代码为我提供了集合中的所有文档。

现在我想对查询应用一些条件。

为此我使用了

config.set("mongo.input.query","{customerId: 'some mongo id'}")

这一次只需要一个条件。如果'usage' > 30,我想添加一个条件

1) 如何使用 spark 和 mongo 为 mongo 查询添加多个条件(包括大于和小于)??

我还想使用 scala 遍历查询结果的每个文档??

2) 如何使用 scala 遍历结果?

【问题讨论】:

  • 这里有一些侧面标志:Mongo 的 hadoop 格式存在资源处理问题,使连接保持打开状态。当我们将它与 Spark 混合时,这是一个爆炸性的组合。 避免
  • @maasg mongo与spark的连接还有其他选择吗??

标签: mongodb scala apache-spark


【解决方案1】:

你好,你可以试试这个:

有一个项目集成了MongoDB和Spark

https://github.com/Stratio/deep-spark/tree/develop

1) 做一个 git clone

2) 进入 deep-spark,然后进入 deep-parent

3) mvn 安装

4) 使用此选项打开 spark-shell:

./spark-shell --jars YOUR_PATH/deep-core-0.7.0-SNAPSHOT.jar,YOUR_PATH/deep-commons-0.7.0-SNAPSHOT.jar,YOUR_PATH/deep-mongodb-0.7.0-SNAPSHOT .jar,YOUR_PATH/mongo-java-driver-2.12.4-sources.jar

记得用真实路径覆盖“YOUR_PATH”

5)在spark shell中执行一个简单的例子:

import com.stratio.deep.mongodb.config.MongoDeepJobConfig
import com.stratio.deep.mongodb.extractor.MongoNativeDBObjectExtractor
import com.stratio.deep.core.context.DeepSparkContext
import com.mongodb.DBObject
import org.apache.spark.rdd.RDD
import com.mongodb.QueryBuilder
import com.mongodb.BasicDBObject

val host = "localhost:27017"


val database = "test"

val inputCollection = "input";

val deepContext: DeepSparkContext = new DeepSparkContext(sc)

val inputConfigEntity: MongoDeepJobConfig[DBObject] = new MongoDeepJobConfig[DBObject](classOf[DBObject])


val query: QueryBuilder  = QueryBuilder.start();

query.and("number").greaterThan(27).lessThan(30);


inputConfigEntity.host(host).database(database).collection(inputCollection).filterQuery(query).setExtractorImplClass(classOf[MongoNativeDBObjectExtractor])


val inputRDDEntity: RDD[DBObject] = deepContext.createRDD(inputConfigEntity)

最好的一点是您可以使用 QueryBuilder 对象进行查询

你也可以像这样传递一个 DBObject:

{ "number" : { "$gt" : 27 , "$lt" : 30}}

如果你想迭代你可以使用方法 yourRDD.collect()。你也可以使用你的RDD.foreach,但你必须提供一个函数。

还有另一种方法可以将罐子添加到 spark 中。您可以修改 spark-env.sh 并将这一行放在最后:

CONFDIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
for jar in $(ls $CONFDIR/../lib/*.jar); do
  SPARK_CLASSPATH=$SPARK_CLASSPATH:${jar}
done

在 lib 文件夹中,您可以放置​​您的库,仅此而已。

免责声明:我目前正在研究 Stratio

【讨论】:

  • 此项目已被弃用,不再有效。这个答案应该被删除。
【解决方案2】:

1) 为了给您的查询添加条件,只需将它们添加到 'mongo.input.query' 提供的字典中:

config.set("mongo.input.query","{customerId: 'some mongo id', usage: {'$gt': 30}")

要更好地了解查询的工作原理,请参阅:

http://docs.mongodb.org/manual/tutorial/query-documents/

http://docs.mongodb.org/getting-started/python/query/

2) 要对结果进行迭代,您可能需要查看触发 RDD 方法“收集”,此链接中的更多信息,只需查找收集方法:

http://spark.apache.org/docs/latest/api/scala/#org.apache.spark.rdd.RDD

【讨论】:

    【解决方案3】:

    是的,你可以使用 mongodb-spark 库,我在这个 Stack Overflow 线程中发布了一些示例:

    How to query MongoDB via Spark for geospatial queries

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-01-20
      • 1970-01-01
      • 1970-01-01
      • 2017-04-17
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多