Spark用dataframe操作ES

1直接上代码: 2package com.suning.scdc.hspark.goods.test 3 4import scala.collection.Seq 5import scala.collection.mutable.LinkedList 6import org.apache.spark.SparkConf 7import org.apache.spark.SparkContext 8import org.elasticsearch.spark.sparkContextFunctions 9import org.slf4j.LoggerFactory 10import org.elasticsearch.spark.rdd.EsSpark 11import org.apache.spark.sql.SparkSession 12import org.apache.spark.sql.SaveMode 13 14object OrderES { 15 val logger = LoggerFactory.getLogger(OrderES.getClass) 16 var sc: SparkContext = null 17 def main(args: Array[String]): Unit = { 18 val conf = new SparkConf() 19 .set("es.nodes", "10.37.154.82,10.37.154.83,10.37.154.84") 20 .set("cluster.name", "elasticsearch") 21 .set("es.port", "9200") 22 23 sc = new SparkContext(conf) 24 dfEs(sc) 25 //esRdd(sc) 26 } 27 28 def esRdd(sc: SparkContext): Unit = { 29 //查询合作方为abc的数据 30 val query = """{"query":{"match":{"memberId": "7013894650"}}}""" 31 val esRdd = sc.esRDD(s"snprime_login/login", query) 32 val rdd = esRdd.map(line => { 33 val key = line._1 34 val value = line._2 35 36 for (tmp <- value) { 37 val key1 = tmp._1 38 val value1 = tmp._2 39 } 40 41 val mp = scala.collection.immutable.Map( 42 "orderNo" -> value("memberId").toString(), 43 "loginTm" -> value("loginTime").toString(), 44 "year" -> "1994") 45 46 (key, mp) 47 48 }) 49 print("lst=") 50 rdd.foreach(println) 51 EsSpark.saveToEsWithMeta(rdd, "bmps/order") 52 } 53 54 55 def dfEs(sc: SparkContext): Unit = { 56 val spark = SparkSession 57 .builder() 58 .appName("sql test") 59 .master("local") 60 .getOrCreate() 61 62 import spark.implicits._ 63 import spark.sql 64 65 // //创建dataframe示例 66 // val df = spark.read.json("C:\\Users\\Administrator\\Desktop\\people.json") 67 // df.createOrReplaceTempView("people") 68 // val sqlDF = spark.sql("select * from people") 69 // sqlDF.show() 70 71 val query = """{"query":{"match":{"memberId": "7013894650"}}}""" 72 73 74 val readDf = spark.read.format("org.elasticsearch.spark.sql").load(s"snprime_login/login") 75 .select("memberId", "loginTime") 76 77 readDf.show 78 79 80 // set primary key for es 81 val esmap = Map("es.mapping.id" -> "memberId") 82 83 readDf.write.format("org.elasticsearch.spark.sql").options(esmap).save("bmps/order") 84 //readDf.write.mode(SaveMode.Append).format("org.elasticsearch.spark.sql").save("bmps/order") 85 86 87 //val esDf = sqlContext.esDF(s"snprime_login/login", query) 88 //将dataFrame/rdd写入es 89 //esRdd.saveToEs("cmall_order/order") 90 //resultDf.saveToEs("index/type") 91 // val schema = StructType( 92 // Seq( 93 // StructField("memberId",StringType,true) 94 // ,StructField("loginTime",StringType,true) 95 // ) 96 // ) 97 // val schema2 = StructType(List( 98 // 99 // StructField("integer_column", IntegerType, nullable = false), 100 // 101 // StructField("string_column", StringType, nullable = true), 102 // 103 // StructField("date_column", DateType, nullable = true) 104 // 105 //)) 106 // 107 108 109 } 110}
点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )