代码如下
private static voidtest2()throwsException {
EnvironmentSettings fsSettings = EnvironmentSettings.newInstance().useOldPlanner().inStreamingMode().build();
StreamExecutionEnvironment fsEnv = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tttt = StreamTableEnvironment.create(fsEnv,fsSettings);
Properties properties =newProperties();
properties.setProperty("bootstrap.servers","localhost:9092");
properties.setProperty("group.id","test"+ System.currentTimeMillis());
FlinkKafkaConsumer011 kafkaConsumer011 =newFlinkKafkaConsumer011("test", newSimpleStringSchema(),properties);
kafkaConsumer011.setStartFromTimestamp(System.currentTimeMillis());
DataStream stream = fsEnv.addSource(kafkaConsumer011);
stream.print();
fsEnv.execute("test");
}2019-10-21 23:06:58.970 [Thread-8] INFO o.a.f.s.connectors.kafka.FlinkKafkaConsumerBase - Consumer subtask 0 creating fetcher with offsets {}.
表现为,无法收到kafka发送过来的消息。这就很奇怪了,因为我使用kafka-consumer-console.sh的方式可以看到实际上是有消息产出的,但是这里死活就收不到数据。