ICode9

精准搜索请尝试: 精确搜索
首页 > 数据库> 文章详细

大数据中各种框架的连接器(Spark, Flink, MongoDB, Kafka, Hive, Hbase等)

2021-05-22 19:03:14  阅读:226  来源: 互联网

标签:val database MongoDB Flink client 连接器 mongodb test TODO


本文不支持复制粘贴的转载,鼓励修改和扩展后的转载。转载前必须邮件取得本人同意,联系方式:1505514388@qq.com

各种语言和框架连接mongodb的基础代码

使用Java客户端连接mongodb

        //获取mongo客户端
        MongoClient mongoClient=new MongoClient("localhost");
        //获取数据库
        MongoDatabase database=mongoClient.getDatabase("test_database");
        //获取集合
        MongoCollection<Document>collection=database.getCollection("test_database");
        //查询集合中的第一个元素
        Document myDoc=collection.find().first();
        System.out.println(myDoc);
        //关闭资源
        mongoClient.close();

使用scala客户端连接mongodb

    val mongoClient = MongoClient("mongodb://localhost:27017")
    val database = mongoClient.getDatabase("test_basedata")
    val collection = database.getCollection("test_basedata")
    val document = collection.find().first()
    System.out.println(document.toString)
    mongoClient.close()

使用spark连接mongodb

在单纯的使用scala的mongodb的连接器时遇到了以下报错:

com.mongodb.ConnectionString.getThreadsAllowedToBlockForConnectionMultiplier()Ljava/lang/Integer;

解决办法为:

<dependency>
    <groupId>org.mongodb</groupId>
    <artifactId>mongo-java-driver</artifactId>
    <version>3.10.0</version>
</dependency>

具体原因未知
spark的连接器的代码如下:

    //TODO 开启环境
    val spark = SparkSession.builder().master("local")
      .config("spark.mongodb.input.uri", "mongodb://127.0.0.1/test_database.test_database")
      .config("spark.mongodb.output.uri", "mongodb://127.0.0.1/test_database.test_database")
      .getOrCreate()
   //TODO 数据操作
    val testDF = MongoSpark.load(spark)
    testDF.show(20)
    //TODO 关闭环境
    spark.close()

使用flink将mongodb作为数据源

在flink中没有将mongodb作为数据源的,所以下面使用的依赖也是第三方连接器。
在一般情况下,也不会遇到将mongodb作为flink的数据源。
所需要的依赖:

<!-- https://mvnrepository.com/artifact/org.mongodb/casbah-core -->
<dependency>
    <groupId>org.mongodb</groupId>
    <artifactId>casbah-core_2.11</artifactId>
    <version>3.1.1</version>
</dependency>

下面是自定义的source

package com.myFlink.test

import com.mongodb.{BasicDBObject, MongoClientURI, casbah}
import com.mongodb.casbah.{MongoClient, MongoClientURI}
import org.apache.flink.configuration.Configuration
import org.apache.flink.streaming.api.functions.source.{RichSourceFunction, SourceFunction}

class MongodbSource extends RichSourceFunction[User]{
  //创建mongodb数据源时运行
  var client:MongoClient=_
  override def open(parameters: Configuration): Unit = {
    //连接本地mongodb
    client= MongoClient(casbah.MongoClientURI("mongodb://localhost:27017"))
  }
  //实时运行的函数
  override def run(sourceContext: SourceFunction.SourceContext[User]): Unit = {
    //TODO 获取数据库
    val database = client("test_database")
    //TODO 获取集合
    val coll = database("test_database")
    //TODO 取出数据
    val query = new BasicDBObject("date", "20171203")
    val cursorType = coll.find(query)
    if(cursorType.nonEmpty){
      val oneData = cursorType.next() //拿出一条数据
      sourceContext.collect(
        User(
          name=oneData.get("name").toString,
          date=oneData.get("date").toString
        )
      )
    }
  }
  //结束时的函数
  override def cancel(): Unit ={
    if(client!=null){
      client.close()
    }
  }
}

case class User(name:String,date:String)

然后向本地的mongdb中插入数据:

db.getCollection("test_database").insert({"name":"kone", "date":"20171203"})

然后输出:

4> User(kone,20171203)

使用flink将mongodb作为sink

package com.myFlink.test
import com.mongodb.BasicDBObject
import com.mongodb.casbah.{MongoClient, MongoClientURI}
import org.apache.flink.configuration.Configuration
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction

class MongoDBSink extends RichSinkFunction[User]{
  var client:MongoClient=_
  override def open(parameters: Configuration): Unit = {
    client=MongoClient(MongoClientURI("mongodb://localhost:27017"))
  }

  override def invoke(value: User): Unit = {
    //TODO 获取数据库
    val database = client("test_database")
    //TODO 获取集合
    val coll = database("test_database")
    //TODO 数据操作
    val obj = new BasicDBObject("name", value.name).append("date", value.date)
    coll.insert(obj)
  }

  override def close(): Unit = {
    if(client!=null){
      client.close()
    }
  }
}

还没完成。

标签:val,database,MongoDB,Flink,client,连接器,mongodb,test,TODO
来源: https://www.cnblogs.com/ALINGMAOMAO/p/14799513.html

本站声明: 1. iCode9 技术分享网(下文简称本站)提供的所有内容,仅供技术学习、探讨和分享;
2. 关于本站的所有留言、评论、转载及引用,纯属内容发起人的个人观点,与本站观点和立场无关;
3. 关于本站的所有言论和文字,纯属内容发起人的个人观点,与本站观点和立场无关;
4. 本站文章均是网友提供,不完全保证技术分享内容的完整性、准确性、时效性、风险性和版权归属;如您发现该文章侵犯了您的权益,可联系我们第一时间进行删除;
5. 本站为非盈利性的个人网站,所有内容不会用来进行牟利,也不会利用任何形式的广告来间接获益,纯粹是为了广大技术爱好者提供技术内容和技术思想的分享性交流网站。

专注分享技术,共同学习,共同进步。侵权联系[81616952@qq.com]

Copyright (C)ICode9.com, All Rights Reserved.

ICode9版权所有