spark寫入ES(動態模板)


使用es-hadoop插件,主要使用elasticsearch-spark-20_2.11-6.2.x.jar

官網:https://www.elastic.co/guide/en/elasticsearch/hadoop/current/reference.html

關於ES詳細的配置參數 大家可以看下面的這個類:

org.elasticsearch.hadoop.cfg.ConfigurationOptions

sparkstreaming寫入ES:
      
 
        SparkConf conf = new SparkConf();
        conf.set("es.index.auto.create", "true");
        conf.set("es.nodes", "10.8.18.16,10.8.18.45,10.8.18.76");
        conf.set("es.port", "9200");
        JavaStreamingContext ssc= null;
        try {
            ssc= new JavaStreamingContext(conf, new Duration(5000L));
            JavaSparkContext jsc =ssc.sparkContext();                        
            String json1 = "{\"reason\" : \"business\",\"airport\" : \"sfo\"}";  
            String json2 = "{\"participants\" : 5,\"airport\" : \"otp\"}";

            JavaRDD<String> stringRDD = jsc.parallelize(ImmutableList.of(json1, json2));
            Queue<JavaRDD<String>> microbatches = new LinkedList<JavaRDD<String>>();      
            microbatches.add(stringRDD);
            JavaDStream<String> stringDStream = ssc.queueStream(microbatches);
            
            //接口1:es的配置通過SparkConf配置
            //使用動態模板,用{}將動態生成的字段名括起來,注意是作用於index
            //而不是type
            //JavaEsSparkStreaming.saveJsonToEs(stringDStream, "spark-{airport}/doc");
            
            Map<String,String> map = new HashMap<String,String>();
            map.put("es.index.auto.create", "true");
            map.put("es.nodes", "ip1,ip2,ip3");
            map.put("es.resource.write", "spark-{airport}/doc");
            map.put("es.port", "9200");
            //接口2:es的配置通過HashMap配置,其中讀取es是index的key為es.resource.read
            //寫入的key為es.resource.write
            //JavaEsSparkStreaming.saveJsonToEs(stringDStream, map);
            //接口3:與接口2類似,只是該接口支持直接填寫index參數
            JavaEsSparkStreaming.saveJsonToEs(stringDStream,"spark-{airport}/doc", map);
            ssc.start();
            ssc.awaitTermination();
        } catch (Throwable e) {
            // TODO 自動生成的 catch 塊
            ssc.close();
            e.printStackTrace();
        }

 

//使用動態模板,用{}將動態生成的字段名括起來,注意是作用於index


免責聲明!

本站轉載的文章為個人學習借鑒使用,本站對版權不負任何法律責任。如果侵犯了您的隱私權益,請聯系本站郵箱yoyou2525@163.com刪除。



 
粵ICP備18138465號   © 2018-2025 CODEPRJ.COM