Advertisement
Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- val realtimeSource: DataStream[CoilMeasurement] = env.addSource(new FlinkKafkaConsumer010[CoilMeasurement](
- realtimeDataKafkaTopic,
- inputSchema,
- loadConsumerKafkaProperties).setStartFromEarliest())
- val flatnessSource: DataStream[CoilMeasurement] = env.addSource(new FlinkKafkaConsumer010[CoilMeasurement](
- flatnessDataKafkaTopic,
- inputSchema,
- loadConsumerKafkaProperties).setStartFromEarliest())
Advertisement
Add Comment
Please, Sign In to add comment
Advertisement