ICode9

精准搜索请尝试: 精确搜索
首页 > 其他分享> 文章详细

练习 : Flink sink to ElasticSearch

2022-07-01 21:03:21  阅读:198  来源: 互联网

标签:flink Flink new streaming ElasticSearch sink org apache import


 

ElasticSearch

package test;

import bean.Stu;
import org.apache.flink.api.common.functions.RuntimeContext;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.source.SourceFunction;
import org.apache.flink.streaming.connectors.elasticsearch.ElasticsearchSinkFunction;
import org.apache.flink.streaming.connectors.elasticsearch.RequestIndexer;
import org.apache.flink.streaming.connectors.elasticsearch7.ElasticsearchSink;
import org.apache.flink.streaming.connectors.redis.common.config.FlinkJedisPoolConfig;
import org.apache.http.HttpHost;
import org.elasticsearch.action.index.IndexRequest;
import org.elasticsearch.client.Requests;

import java.util.*;

public class SinkToElasticsearch {
    public static void main(String[] args) {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);
        DataStreamSource<Stu> source = env.addSource(new SourceFunction<Stu>() {
            private boolean running = true;

            @Override
            public void run(SourceContext<Stu> sourceContext) throws Exception {
                while (running) {
                    for (int i = 0; i < 10; i++) {
                        ArrayList<String> subs = new ArrayList<String>(Arrays.asList("语文", "数学", "英语", "化学", "物理", "生物"));
                        List<String> names = Arrays.asList("张三", "李四", "王五", "赵六","田七");
                        int next = new Random().nextInt(15);
                        int random = new Random().nextInt(101);
                        sourceContext.collect(new Stu(names.get(next * random % 5), subs.get(next * random % 6), random));
                        Thread.sleep(1000);
                    }
                    Thread.sleep(20000);
                }
            }

            @Override
            public void cancel() {
                running = false;
            }
        });

       //创建一个 hosts 列表
        ArrayList<HttpHost> httpHosts = new ArrayList<>();
        httpHosts.add(new HttpHost("hadoop106",9200));
        FlinkJedisPoolConfig config = new FlinkJedisPoolConfig.Builder()
                .setHost("hadoop106")
                .build();

        // elasticsearch  Sink Function
        ElasticsearchSinkFunction<Stu> elasticsearchSinkFunction = new ElasticsearchSinkFunction<Stu>() {
            @Override
            public void process(Stu stu, RuntimeContext ctx, RequestIndexer indexer) {
                HashMap<String, String> map = new HashMap<>();
                map.put(stu.getName(), stu.getSub() + " " + stu.getScore());

                IndexRequest request = Requests.indexRequest()
                        .index("clicks")
                        .type("type")
                        .source(map);

                indexer.add(request);

            }
        };

        //写入 es
        source.addSink(new ElasticsearchSink.Builder<>(httpHosts,elasticsearchSinkFunction).build());
        try {
            env.execute();
        } catch (Exception e) {
            e.printStackTrace();
        }
        
        }


}

 

标签:flink,Flink,new,streaming,ElasticSearch,sink,org,apache,import
来源: https://www.cnblogs.com/chang09/p/16435950.html

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

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

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

ICode9版权所有