KafkaSource source = KafkaSource.builder() .setBootstrapServers("localhost:9092") .setTopics("source") .setGroupId("my-group") .setStartingOffsets(OffsetsInitializer.latest()) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); DataStream sourceStream = env.fromSource( source, WatermarkStrategy.forMonotonousTimestamps(), "Kafka Source") .uid("kafkasourceuid");