【发布时间】:2020-09-18 09:56:00
【问题描述】:
我尝试提交简单的 flink 作业以接受来自 kafka 的消息,但在提交作业后不到一分钟,作业失败并出现以下 kafka 异常。我在本地机器上运行了 kafka 2.12,并且我已经配置了该作业使用的主题。
public static void main(String[] args) throws Exception {
Properties properties = new Properties();
properties.setProperty("bootstrap.servers", "127.0.0.1:9092");
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<String> kafkaData = env
.addSource(new FlinkKafkaConsumer<String>("test-topic",
new SimpleStringSchema(), properties));
kafkaData.print();
env.execute("Aggregation Job");
}
这是一个例外:
Job has been submitted with JobID 5cc30fe72f685406126e2f5a26f10341
------------------------------------------------------------
The program finished with the following exception:
org.apache.flink.client.program.ProgramInvocationException: The main method caused an error: org.apache.flink.client.program.ProgramInvocationException: Job failed (JobID: 5cc30fe72f685406126e2f5a26f10341)
at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:335)
...
Caused by: org.apache.kafka.common.errors.TimeoutException: Timeout expired while fetching topic metadata
I saw another question in stackoverflow,但这并不能解决问题。我没有在 kafka 代理上配置任何 SSL。任何建议将不胜感激。
【问题讨论】: