Loading src/main/java/com/censoft/flink/StreamingJob.java +2 −2 Original line number Diff line number Diff line Loading @@ -54,7 +54,7 @@ public class StreamingJob { //3、4 Object -> String streamOperator = streamOperator.map(algorithmPushDto ->{ System.out.println(algorithmPushDto); System.out.println(algorithmPushDto.toString()); return algorithmPushDto; }); Loading @@ -80,7 +80,7 @@ public class StreamingJob { SingleOutputStreamOperator<String> outputStreamOperator = streamOperator.map(AlgorithmPushDto::toString); //3、5 输出kafka outputStreamOperator.addSink(new FlinkKafkaProducer("172.16.33.152:9092", "test-topic", new SimpleStringSchema())); // outputStreamOperator.addSink(new FlinkKafkaProducer("172.16.33.152:9092", "test-topic", new SimpleStringSchema())); //3、6 打印输出 outputStreamOperator.print(); Loading src/main/java/com/censoft/flink/mqtt/MqttConsumer.java +1 −0 Original line number Diff line number Diff line Loading @@ -104,6 +104,7 @@ public class MqttConsumer extends RichParallelSourceFunction<AlgorithmPushDto> { public void messageArrived(String s, MqttMessage message) throws Exception { //订阅消息字符 String msg = new String(message.getPayload()); System.out.println(msg); byte[] bymsg = getBytesFromObject(msg); AlgorithmPushDto algorithmPushDto = JSON.parseObject(msg, AlgorithmPushDto.class); Loading Loading
src/main/java/com/censoft/flink/StreamingJob.java +2 −2 Original line number Diff line number Diff line Loading @@ -54,7 +54,7 @@ public class StreamingJob { //3、4 Object -> String streamOperator = streamOperator.map(algorithmPushDto ->{ System.out.println(algorithmPushDto); System.out.println(algorithmPushDto.toString()); return algorithmPushDto; }); Loading @@ -80,7 +80,7 @@ public class StreamingJob { SingleOutputStreamOperator<String> outputStreamOperator = streamOperator.map(AlgorithmPushDto::toString); //3、5 输出kafka outputStreamOperator.addSink(new FlinkKafkaProducer("172.16.33.152:9092", "test-topic", new SimpleStringSchema())); // outputStreamOperator.addSink(new FlinkKafkaProducer("172.16.33.152:9092", "test-topic", new SimpleStringSchema())); //3、6 打印输出 outputStreamOperator.print(); Loading
src/main/java/com/censoft/flink/mqtt/MqttConsumer.java +1 −0 Original line number Diff line number Diff line Loading @@ -104,6 +104,7 @@ public class MqttConsumer extends RichParallelSourceFunction<AlgorithmPushDto> { public void messageArrived(String s, MqttMessage message) throws Exception { //订阅消息字符 String msg = new String(message.getPayload()); System.out.println(msg); byte[] bymsg = getBytesFromObject(msg); AlgorithmPushDto algorithmPushDto = JSON.parseObject(msg, AlgorithmPushDto.class); Loading