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}
Spark用dataframe操作ES
Stella981
2021-10-12
1205 0 0
点赞
收藏
评论区
加载中...