【问题标题】:Flink, Kafka and Zookeeper with an URI带有 URI 的 Flink、Kafka 和 Zookeeper
【发布时间】:2016-11-29 20:38:01
【问题描述】:

我正在尝试从本地计算机连接到 Kafka:

kafkaParams.setProperty("bootstrap.servers", Defaults.BROKER_URL)
kafkaParams.setProperty("metadata.broker.list", Defaults.BROKER_URL)
kafkaParams.setProperty("group.id", "group_id")
kafkaParams.setProperty("auto.offset.reset", "earliest")

完全没问题,但是我的BROKER_URI 定义如下my-server.com:1234/my/subdirectory

我发现这种现象称为 chroot 路径。

它会抛出以下错误:Caused by: org.apache.kafka.common.config.ConfigException: Invalid url in bootstrap.servers: my-server.com:1234/my/subdirectory

我该如何解决这个问题?

这些是我的依赖项:

val flinkVersion = "1.0.3"

"org.apache.flink" %% "flink-scala" % flinkVersion % "provided",
"org.apache.flink" %% "flink-streaming-scala" % flinkVersion % "provided",
"org.apache.flink" %% "flink-connector-kafka-0.9" % flinkVersion,

【问题讨论】:

    标签: scala hadoop apache-kafka kafka-consumer-api apache-flink


    【解决方案1】:

    只需尝试host:port 格式,不要使用路径上下文和斜线。如果您有多个服务器,它将是一个列表host1:port1,host2:port2

    参考:http://kafka.apache.org/documentation.html

    【讨论】:

    • 这给出了以下错误:Exception in thread "main" org.apache.kafka.common.errors.TimeoutException: Timeout expired while fetching topic metadata
    • 表示配置格式没问题。接下来要寻找的是 kafka 实例上是否运行 iptables 或防火墙。你能从你的客户端远程登录 kafka 实例吗?
    • 有趣的是,我可以使用 Kafka 控制台消费者进行连接:./kafka-console-consumer.sh --zookeeper my-server.com:1234/my/subdirectory --topic my-topic --from-beginning 工作得非常好。 Telnet 也可以正常工作:telnet my-server.com 1234
    【解决方案2】:

    bootstrap.servers 应该是一个逗号分隔的列表,如下所示:address1:port1,address2:port2,...,addressn:portn。如果您只有一个 Kafka 代理,您应该输入类似 localhost:9092 的内容(除非您将 Kafka 配置为在另一个端口上运行)。

    您可以参考this post from dataArtisans 了解更多关于如何使 Flink 和 Kafka 协同工作的详细信息。

    【讨论】:

      【解决方案3】:

      愚蠢。动物园管理员!=卡夫卡。正如您在代码中看到的那样,我两次使用了相同的 URL,但结果它们应该是不同的。

      我正在尝试从本地计算机连接到 Kafka:

      kafkaParams.setProperty("bootstrap.servers", Defaults.KAFKA_URL)
      kafkaParams.setProperty("metadata.broker.list", Defaults.ZOOKEEPER_URL)
      kafkaParams.setProperty("group.id", "group_id")
      kafkaParams.setProperty("auto.offset.reset", "earliest")
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2018-07-07
        • 1970-01-01
        • 2019-02-08
        • 1970-01-01
        • 2017-11-20
        • 2018-06-17
        • 2017-05-11
        • 2022-01-19
        相关资源
        最近更新 更多