Flink(7) 自定義數(shù)據(jù)源

簡介

只要實現(xiàn) SourceFunction 接口對應的方法就可以自定義數(shù)據(jù)源
1.創(chuàng)建環(huán)境

 public static void main(String[] args) throws Exception {

        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        DataStreamSource<SensorReading> streamSource = env.addSource(new MySensorSource());

        streamSource.print();

        env.execute();

    }

2.實現(xiàn) SourceFunction 接口

 public static class MySensorSource implements SourceFunction<SensorReading> {

        //定義一個標識位用來控制數(shù)據(jù)
        private boolean running = true;
        //定義一個隨機數(shù)發(fā)生器
        Random random = new Random();

        public void run(SourceContext<SensorReading> ctx) throws Exception {

            //設(shè)置10個傳感器的初識溫度
            HashMap<String, Double> sensorTempMap = new HashMap<String, Double>();
            for (int i = 0; i < 10; i++) {
                sensorTempMap.put("sensor_" + (i + 1), 60 + random.nextGaussian() * 20);
            }

            while (running) {

                for(String sensorId : sensorTempMap.keySet()){
                    //在當前的溫度基礎(chǔ)上隨機波動
                    Double newtemp = sensorTempMap.get(sensorId) + random.nextGaussian();

                    sensorTempMap.put(sensorId,newtemp);

                    ctx.collect(new SensorReading(sensorId,System.currentTimeMillis(),newtemp));
                }
                //控制輸出頻率
                Thread.sleep(1000L);
            }

        }

        public void cancel() {
            running = false;
        }
    }
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請聯(lián)系作者
【社區(qū)內(nèi)容提示】社區(qū)部分內(nèi)容疑似由AI輔助生成,瀏覽時請結(jié)合常識與多方信息審慎甄別。
平臺聲明:文章內(nèi)容(如有圖片或視頻亦包括在內(nèi))由作者上傳并發(fā)布,文章內(nèi)容僅代表作者本人觀點,簡書系信息發(fā)布平臺,僅提供信息存儲服務。

相關(guān)閱讀更多精彩內(nèi)容

友情鏈接更多精彩內(nèi)容