【问题标题】:SparkR 2.x application using MongoDB and RStudio使用 MongoDB 和 RStudio 的 SparkR 2.x 应用程序
【发布时间】:2017-02-26 08:01:32
【问题描述】:

我正在尝试锻炼一个 Apache Spark 应用程序,该应用程序应该在 MongoDB 数据库上运行聚合查询并写回结果。我能够解决问题的 Java 版本,但现在需要使用 RStudio 将其移植到 R 语言。

工作的Java版本:-

public static void main(String args[]) {

SparkConf sparkConf = new SparkConf(true)
        .setMaster("local[*]")
        .setSparkHome(SPARK_HOME)
        .setAppName("SparklingMongoApp")
        .set("spark.ui.enabled", "false")
        .set("spark.app.id", APP)
        .set("spark.mongodb.input.uri", "mongodb://admin:password@host:27017/input_collection")
        .set("spark.mongodb.output.uri", "mongodb://admin:password@host:27017/output_collection");


JavaSparkContext javaSparkContext = new JavaSparkContext(sparkConf);
JavaMongoRDD<Document> javaMongoRDD = MongoSpark.load(javaSparkContext);

Dataset<Row> dataset = javaMongoRDD.toDF();

dataset.createOrReplaceTempView(TEMP_VIEW);

// a valid spark sql QUERY
Dataset<Row> computedDataSet = dataset.sqlContext().sql(QUERY);
MongoSpark.save(computedDataSet);
javaSparkContext.close();

}

我正在尝试锻炼的等效 R/RStudio 版本:-

library(SparkR, lib.loc = c(file.path(Sys.getenv("SPARK_HOME"), "R", "lib")))

##PROBLEM - Is this correct way of setting configuration?
sparkConfig <- list("spark.driver.memory"="1g","spark.mongodb.input.uri"="mongodb://username:password@localhost:27017/price_subset?authSource=admin","spark.mongodb.output.uri"="mongodb://username:password@localhost:27017/price_subset_output?authSource=admin")

customSparkPackages <- c("org.mongodb.spark:mongo-spark1-connector_2.11:1.0.0");

##Starting Up: SparkSession
##PROBLEM-1 Is this correct way of initializing spark session ?
sparkSession <- sparkR.session(appName="MongoSparkConnectorTour",master = "local[*]",enableHiveSupport = FALSE,sparkConfig = sparkConfig,sparkPackages = customSparkPackages)


##PROBLEM-2 - This complains about being deprecated. How to fix this ?
sqlContext <- sparkRSQL.init(sparkSession)

## Save some data
charactersRdf <- data.frame(list(name=c("Bilbo Baggins", "Gandalf", "Thorin", "Balin", "Kili", "Dwalin", "Oin", "Gloin", "Fili", "Bombur"),
                                 age=c(50, 1000, 195, 178, 77, 169, 167, 158, 82, NA)))

charactersSparkdf <- createDataFrame(sqlContext, charactersRdf)
#PROBLEM-3 This throws an error - Error in invokeJava(isStatic = FALSE, objId$id, methodName, ...) : 
#  java.lang.NoClassDefFoundError: com/mongodb/ConnectionString
write.df(charactersSparkdf, "", source = "com.mongodb.spark.sql.DefaultSource", mode = "overwrite")

我尝试关注 SparkR 文档,但仍然无法锻炼运行示例。

期望:-

  1. 在 RStudio 中初始化 spark 会话的正确方法是什么。 MongoDB official sample 不适用于我,因为它仅适用于 SparkShell(它会挂在我的机器上)并且已弃用。我想要可以在 RStudio 中运行的代码 sn-p。

  2. 如何修复 java.lang.NoClassDefFoundError。

任何 SparkR 2.x + MongoDB 3.x 代码的示例/参考都将受到高度赞赏。

版本:- 阿帕奇星火 - 2.0.1 爪哇 - 1.8 MongoDB - 3 R - 最新的

【问题讨论】:

  • 如果你正在使用 RStudio,你也可以试试他们的 SparklyR 包(它与 RStudio 预览版中的 IDE 集成)。 spark.rstudio.com

标签: r mongodb apache-spark rstudio sparkr


【解决方案1】:

终于搞定了。原来 MongoDB 文档有 Spark 1.6 的示例,而我运行的是 Spark 2.0.1。

无论如何,这就是使用 RStudio 对我有用的方法:-

 ## Make sure you have SPARK_HOME environment variable set to your spark home director.
library(SparkR, lib.loc = c(file.path(Sys.getenv("SPARK_HOME"), "R", "lib")))

spark <- sparkR.session(master="local[*]", appName = "mongoSparkR",enableHiveSupport = FALSE,sparkPackages = c("org.mongodb.spark:mongo-spark-connector_2.11:2.0.0-rc0"),sparkConfig = list(spark.mongodb.input.uri="mongodb://username:password@hostname:27017/database.collection_name?authSource=admin",spark.mongodb.output.uri="mongodb://username:password@hostname:27017/database.collection_name_output?authSource=admin"))

pricing_df <- read.df(source = "com.mongodb.spark.sql.DefaultSource",x=10000)
head(pricing_df)
createOrReplaceTempView(pricing_df,"T_YOUR_TABLE")

 ## Obviously this is just a dummy SQL, replace with it yours.
result_df <- sql("SELECT year(price) as YEAR, month(price) as MONTH , SUM(midPrice) as SUM_PRICING_DATA FROM T_YOUR_TABLE GROUP BY year(price),month(price)  ORDER BY year(price),month(price)")


 ## stop instance when done.
sparkR.stop()

确保您的 SPARK_HOME/jars 文件夹中有相关的 jar。

我放置的额外 jars(版本可能会随着时间的推移而演变)以使其正常工作:-

org.mongodb.spark_mongo-spark-connector_2.11-2.0.0-rc0.jar

org.mongodb_mongo-java-driver-3.2.2.jar

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2016-09-20
    • 1970-01-01
    • 2020-08-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-09-21
    • 1970-01-01
    相关资源
    最近更新 更多