Authored by gemingdan

EventTime

@@ -20,7 +20,7 @@ public class TraceFlinkExecutor { @@ -20,7 +20,7 @@ public class TraceFlinkExecutor {
20 20
21 public static void main(String[] args) throws Exception{ 21 public static void main(String[] args) throws Exception{
22 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); 22 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
23 - env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime); 23 + env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
24 24
25 FlinkKafkaConsumerBase<String> kafkaConsumer = KafkaUtils.finkKafkaConsumer(PropertiesFactory.kafka().getBrokers(), PropertiesFactory.kafka().getGroup(), PropertiesFactory.kafka().getTopic()); 25 FlinkKafkaConsumerBase<String> kafkaConsumer = KafkaUtils.finkKafkaConsumer(PropertiesFactory.kafka().getBrokers(), PropertiesFactory.kafka().getGroup(), PropertiesFactory.kafka().getTopic());
26 SingleOutputStreamOperator<Spans> clickStream = env.addSource(kafkaConsumer) 26 SingleOutputStreamOperator<Spans> clickStream = env.addSource(kafkaConsumer)