【问题标题】:Flink with kafka issue: Timeout expired while fetching topic metadataFlink 与 kafka 问题:获取主题元数据时超时
【发布时间】: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。任何建议将不胜感激。

【问题讨论】:

    标签: apache-kafka apache-flink


    【解决方案1】:

    我今天也遇到了同样的问题。就我而言,问题在于我未能将我的 flink 应用程序放入 VPC(我的 MSK 集群位于 VPC 中)。编辑 flink 应用程序并将其移动到适当的 VPC 后,问题就消失了。

    我意识到这个问题已经有几个月的历史了,但我想我会发布我的发现,以防其他人碰巧像我一样从 Google 搜索中看到这个问题。

    【讨论】:

      猜你喜欢
      • 2021-10-07
      • 2019-06-12
      • 2019-11-22
      • 2021-03-19
      • 1970-01-01
      • 2019-12-25
      • 1970-01-01
      • 2021-01-05
      • 2020-08-07
      相关资源
      最近更新 更多