【问题标题】:Can spark sql create nosql table?spark sql可以创建nosql表吗?
【发布时间】:2020-07-23 07:21:53
【问题描述】:

到目前为止,我认为SparkSQL可以创建hive表,因为spark-sql也是在shark的基础上改进的,shark from hive。

hive create table

spark-sql 是否创建其他 nosql 表?例如 hbase、cassandra、elasticsearch。

查了相关文档,没有找到spark-sql create table api, 以后会支持吗?

【问题讨论】:

    标签: apache-spark apache-spark-sql nosql


    【解决方案1】:

    目前无法直接通过 spark-sql 选项在 Elasticsearch、Hbase 和 Cassandra 中创建索引。 您可以使用 hbase-client 选项与 hbase 进行本机交互。 - https://mvnrepository.com/artifact/org.apache.hbase/hbase-client

    1.在 Hbase 中创建表

          import org.apache.hadoop.hbase.*;
          import org.apache.hadoop.conf.Configuration;
          val admin = new HBaseAdmin(conf)
    
          if (!admin.tableExists(myTable)) {
          val htd = new HTableDescriptor(myTable)
          htd.addFamily(new HColumnDescriptor("id"))
          htd.addFamily(new HColumnDescriptor("name"))
          htd.addFamily(new HColumnDescriptor("country))
          htd.addFamily(new HColumnDescriptor("pincode"))
          admin.createTable(htd)
    }
    

    2。写入已创建的 Hbase 表。

    import org.apache.spark.streaming.kafka010.KafkaUtils
    import org.apache.kafka.common.serialization.StringDeserializer
    import org.apache.spark.{SparkConf, SparkContext}
    import org.apache.spark.streaming.{Seconds, StreamingContext}
    import org.apache.spark.streaming.kafka010.LocationStrategies.PreferConsistent
    import org.apache.spark.streaming.kafka010.ConsumerStrategies.Subscribe
    import org.apache.hadoop.hbase.{HBaseConfiguration, HColumnDescriptor, HTableDescriptor}
    import org.apache.spark.sql.SparkSession
    import org.apache.hadoop.hbase.client.ConnectionFactory
    import org.apache.hadoop.hbase.client.Connection
    import org.apache.hadoop.hbase.client.Put
    import org.apache.hadoop.hbase.util.Bytes
    
            val kafkaParams = Map[String, Object](
                "bootstrap.servers" -> "cmaster.localcloud.com:9092,cworker2.localcloud.com:9092,cworker1.localcloud.com:9092",
                "key.deserializer" -> classOf[StringDeserializer],
                "value.deserializer" -> classOf[StringDeserializer],
                "group.id" -> "use_a_separate_group_id_for_each_stream",
                "auto.offset.reset" -> "latest",
                "enable.auto.commit" -> (false: java.lang.Boolean)
    
            )
    
            val topics = Array("test")
            val spark = SparkSession.builder().master("local[8]").appName("KafkaSparkHBasePipeline").getOrCreate()
            spark.sparkContext.setLogLevel("OFF")
            val result = kafkaStream.map(record => (record.key, record.value))
            result.foreachRDD(x => {
                val cols = x.map(x => x._2.split(","))
                val arr =cols.map(x => {
                    val id = x(0)
                    val name = x(1)
                    val country = x(2)
                    val pincode = x(3)
                    (id,name,country,pincode)
                })
                arr.foreachPartition { iter =>
                    val conf = HBaseConfiguration.create()
                    conf.set("hbase.zookeeper.quorum", "cmaster.localcloud.com")
                    conf.set("hbase.rootdir", "hdfs://localhost:8020/hbase")
                    conf.set("hbase.zookeeper.property.clientPort", "2181")
                    conf.set("zookeeper.znode.parent", "/hbase")
                    conf.set("hbase.unsafe.stream.capability.enforce", "false")
                    conf.set("hbase.cluster.distributed", "true")
                    val conn = ConnectionFactory.createConnection(conf)
                    import org.apache.hadoop.hbase.TableName
                    val tableName = "htd"
                    val table = TableName.valueOf(tableName)
                    val HbaseTable = conn.getTable(table)
                    val cfPersonal = "personal"
                    iter.foreach(x => {
                        val keyValue = "Key_" + x._1
                        val id = new Put(Bytes.toBytes(keyValue))
                        val name = x._2.toString
                        val country = x._3.toString
                        val pincode = x._4.toString
                        id.addColumn(Bytes.toBytes(cfPersonal), Bytes.toBytes("name"), Bytes.toBytes(name))
                        id.addColumn(Bytes.toBytes(cfPersonal), Bytes.toBytes("country"), Bytes.toBytes(country))
                        id.addColumn(Bytes.toBytes(cfPersonal), Bytes.toBytes("pincode"), Bytes.toBytes(pincode))
                        HbaseTable.put(id)
                    })
                    HbaseTable.close()
                    conn.close()
                }
            })
    

    3.我在这个项目中使用的依赖项,

    name := "KafkaSparkHBasePipeline"
    version := "0.1"
    scalaVersion := "2.11.8"
    resolvers += "Mavenrepository" at "https://mvnrepository.com"
    libraryDependencies += "org.apache.spark" %% "spark-core" % "2.4.2"
    libraryDependencies += "org.apache.spark" %% "spark-sql" % "2.4.2"
    libraryDependencies += "org.apache.hbase" % "hbase-client" % "2.2.0"
    libraryDependencies += "org.apache.kafka" %% "kafka" % "2.3.0"
    libraryDependencies += "org.apache.spark" %% "spark-streaming" % "2.4.3"
    libraryDependencies += "org.apache.spark" %% "spark-streaming-kafka-0-10" % "2.4.3"
    

    4.验证 hbase 中的数据。

    smart@cmaster sathyadev]$ hbase shell 
    Java HotSpot(TM) 64-Bit Server VM warning: Using incremental CMS is deprecated and will likely be removed in a future release HBase Shell Use "help" to get list of supported commands. Use "exit" to quit this interactive shell. 
    For Reference, please visit: http://hbase.apache.org/2.0/book.html
    #shell Version 2.1.0-cdh6.2.1, rUnknown, Wed Sep 11 01:05:56 PDT 2019 Took 0.0064 seconds 
    hbase(main):001:0> scan 'htd'; 
    ROW COLUMN+CELL Key_1010 column=personal:country, timestamp=1584125319078, value=USA Key_1010 column=personal:name, timestamp=1584125319078, value=Mark Key_1010 column=personal:pincode, timestamp=1584125319078, value=54321 Key_1011 column=personal:country, timestamp=1584125320073, value=CA
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-11-26
      • 2020-06-15
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多