代码之家  ›  专栏  ›  技术社区  ›  Daniil Andreyevich Baunov

Spark Mongodb连接器Scala-缺少数据库名称

  •  8
  • Daniil Andreyevich Baunov  · 技术社区  · 8 年前

    我遇到了一个奇怪的问题。我正在尝试使用MongoDB Spark连接器将Spark本地连接到MongoDB。

    val readConfig = ReadConfig(Map("uri" -> "mongodb://localhost:27017/movie_db.movie_ratings", "readPreference.name" -> "secondaryPreferred"), Some(ReadConfig(sc)))
    val writeConfig = WriteConfig(Map("uri" -> "mongodb://127.0.0.1/movie_db.movie_ratings"))
    
    // Load the movie rating data from Mongo DB
    val movieRatings = MongoSpark.load(sc, readConfig).toDF()
    
    movieRatings.show(100)
    

    java.lang.IllegalArgumentException: Missing database name. Set via the 'spark.mongodb.input.uri' or 'spark.mongodb.input.database' property.
    

    在线设置readConfig。我不明白为什么它抱怨我在地图中有一个uri属性,却没有设置uri。 我可能错过了什么。

    2 回复  |  直到 8 年前
        1
  •  11
  •   mrsrinivas    8 年前

    你可以从 SparkSession

    val spark = SparkSession.builder()
        .master("local")
        .appName("MongoSparkConnectorIntro")
        .config("spark.mongodb.input.uri", "mongodb://localhost:27017/movie_db.movie_ratings")
        .config("spark.mongodb.input.readPreference.name", "secondaryPreferred")
        .config("spark.mongodb.output.uri", "mongodb://127.0.0.1/movie_db.movie_ratings")
        .getOrCreate()
    

    使用配置创建数据帧

    val readConfig = ReadConfig(Map("uri" -> "mongodb://localhost:27017/movie_db.movie_ratings", "readPreference.name" -> "secondaryPreferred"))
    val df = MongoSpark.load(spark)
    

    将df写入mongodb

    MongoSpark.save(
    df.write
        .option("spark.mongodb.output.uri", "mongodb://127.0.0.1/movie_db.movie_ratings")
        .mode("overwrite"))
    

    在您的代码中:

    val readConfig = ReadConfig(Map(
        "spark.mongodb.input.uri" -> "mongodb://localhost:27017/movie_db.movie_ratings", 
        "spark.mongodb.input.readPreference.name" -> "secondaryPreferred"), 
        Some(ReadConfig(sc)))
    
    val writeConfig = WriteConfig(Map(
        "spark.mongodb.output.uri" -> "mongodb://127.0.0.1/movie_db.movie_ratings"))
    
        2
  •  1
  •   Vignesh Jaganathan    7 年前

    SparkSession sparkSession = SparkSession.builder()
        .master("local")
        .appName("MongoSparkConnector")
        .config("spark.mongodb.input.uri","mongodb://localhost:27017/movie_db.movie_ratings")
        .config("spark.mongodb.input.readPreference.name", "secondaryPreferred")
        .config("spark.mongodb.output.uri", "mongodb://127.0.0.1/movie_db.movie_ratings")
        .getOrCreate()
    

     SparkSession sparkSession = SparkSession.builder()
            .master("local")
            .appName("MongoSparkConnector")
            .getOrCreate()
    

    然后

         String mongoUrl = "mongodb://localhost:27017/movie_db.movie_ratings";
       sparkSession.sparkContext().conf().set("spark.mongodb.input.uri", mongoURL);
       sparkSession.sparkContext().conf().set("spark.mongodb.output.uri", mongoURL);
       Map<String, String> readOverrides = new HashMap<String, String>();
       readOverrides.put("collection", sourceCollection);
       readOverrides.put("readPreference.name", "secondaryPreferred");
       ReadConfig readConfig = ReadConfig.create(sparkSession).withOptions(readOverrides);
       Dataset<Row> df = MongoSpark.loadAndInferSchema(sparkSession,readConfig);
    
    推荐文章