Ro66ie_www

Untitled

Oct 2nd, 2017
133
0
Never
Not a member of Pastebin yet? Sign Up, it unlocks many cool features!
Java 0.62 KB | None | 0 0
  1.         final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  2.         env.setParallelism(5);
  3.        
  4.         ParameterTool parameterTool = ParameterTool.fromArgs(args);
  5.  
  6.         DataStream<String> messageStream = env.addSource(new FlinkKafkaConsumer010<>(parameterTool.getRequired("topic"),
  7.             new SimpleStringSchema(), parameterTool.getProperties()));
  8.  
  9.         messageStream.map(new MapFunction<String, String>() {
  10.             @Override
  11.             public String map(String value) throws Exception {
  12.                 //--- do something --//
  13.            
  14.                 return "KafkaJob output: " + outputBuilder.toString();
  15.             }
  16.         })
  17.         .rebalance().print();
Advertisement
Add Comment
Please, Sign In to add comment