【问题标题】:Only one node of a Kafka Connect cluster responds to REST API requestsKafka Connect 集群只有一个节点响应 REST API 请求
【发布时间】:2020-08-08 09:59:15
【问题描述】:

我在不同的主机上运行 Kafka Connect 集群,并且我看到这些节点中只有一个响应 REST API 请求(特别是:POST、PUT 或 DELETE 请求)的行为。我可以通过关闭一个节点并向另一个活动节点发出写入命令来可靠地交换响应 API 请求的节点。

这是我的 docker-compose worker 配置:

version: '2'
services:
  connect:
    image: debezium/connect:1.1.0.Final
    ports:
     - 8083:8083
    volumes:
     - /etc/kafka/secrets:/etc/kafka/secrets
    environment:
     - BOOTSTRAP_SERVERS=my.region.aws.confluent.cloud:9092
     - GROUP_ID=debezium-postgres
     - CONFIG_STORAGE_TOPIC=dbz_pg_connect_configs
     - OFFSET_STORAGE_TOPIC=dbz_pg_connect_offsets
     - STATUS_STORAGE_TOPIC=dbz_pg_connect_statuses
     - CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR=3
     - CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR=3
     - CONNECT_STATUS_STORAGE_REPLICATION_FACTOR=3
     - OFFSET_FLUSH_INTERVAL_MS=8000
     - OFFSET_FLUSH_TIMEOUT_MS=60000
     - CONNECT_SECURITY_PROTOCOL=SASL_SSL
     - CONNECT_SASL_MECHANISM=PLAIN
     - CONNECT_SASL_JAAS_CONFIG=org.apache.kafka.common.security.plain.PlainLoginModule required username="<user>" password="<pass>";
     - CONNECT_PRODUCER_SECURITY_PROTOCOL=SASL_SSL
     - CONNECT_PRODUCER_SASL_JAAS_CONFIG=org.apache.kafka.common.security.plain.PlainLoginModule required username="<user>" password="<pass>";
     - CONNECT_PRODUCER_SASL_MECHANISM=PLAIN
     - CONNECT_CONSUMER_SECURITY_PROTOCOL=SASL_SSL
     - CONNECT_CONSUMER_SASL_JAAS_CONFIG=org.apache.kafka.common.security.plain.PlainLoginModule required username="<user>" password="<pass>";
     - CONNECT_CONSUMER_SASL_MECHANISM=PLAIN

我可以使用 Debezium Postgres 连接器和 Kafka Snowflake 连接器来重现这一点。所以我认为问题在于 Kafka Connect REST API 本身,而不是任何特定的连接器库。

根据docs

默认情况下,此服务在端口 8083 上运行。在分布式模式下执行时,REST API 将成为集群的主要接口。 您可以向任何集群成员发出请求;如果需要,REST API 会自动转发请求。

这是我的设置:

  • 2 个唯一的主机
  • 2个运行Kafka Connect的dockerized容器(镜像实际上是debezium/connect:1.1.0.Final)
  • 两者都在端口 8083 上运行 REST 服务。文档中没有表明当容器位于不同主机上时这是一个问题

我看到的行为是这样的:

  • GET 请求在任何情况下都适用于两个节点
  • POST/PUT/DELETE 请求在第一个节点上工作以接受其中一个调用。之后,只有该节点响应 POST/PUT/DELETE。

另一个节点响应:

HTTP/1.1 100 Continue

HTTP/1.1 500 Internal Server Error
Date: Fri, 24 Apr 2020 17:59:55 GMT
Content-Type: application/json
Content-Length: 120
Server: Jetty(9.4.20.v20190813)

{"error_code":500,"message":"IO Error trying to forward REST request: java.net.SocketTimeoutException: Connect Timeout"}

编辑:这里是 Kafka Connect 日志:

Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m 2020-04-24 19:11:14,934 ERROR  ||  IO error forwarding REST request:    [org.apache.kafka.connect.runtime.rest.RestClient]
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m java.util.concurrent.ExecutionException: java.net.SocketTimeoutException: Connect Timeout
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.client.util.FutureResponseListener.getResult(FutureResponseListener.java:118)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.client.util.FutureResponseListener.get(FutureResponseListener.java:101)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.client.HttpRequest.send(HttpRequest.java:685)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.apache.kafka.connect.runtime.rest.RestClient.httpRequest(RestClient.java:125)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.apache.kafka.connect.runtime.rest.RestClient.httpRequest(RestClient.java:65)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.apache.kafka.connect.runtime.rest.resources.ConnectorsResource.completeOrForwardRequest(ConnectorsResource.java:315)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.apache.kafka.connect.runtime.rest.resources.ConnectorsResource.createConnector(ConnectorsResource.java:143)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at java.base/java.lang.reflect.Method.invoke(Method.java:566)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.server.model.internal.ResourceMethodInvocationHandlerFactory.lambda$static$0(ResourceMethodInvocationHandlerFactory.java:52)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.server.model.internal.AbstractJavaResourceMethodDispatcher$1.run(AbstractJavaResourceMethodDispatcher.java:124)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.server.model.internal.AbstractJavaResourceMethodDispatcher.invoke(AbstractJavaResourceMethodDispatcher.java:167)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.server.model.internal.JavaResourceMethodDispatcherProvider$ResponseOutInvoker.doDispatch(JavaResourceMethodDispatcherProvider.java:176)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.server.model.internal.AbstractJavaResourceMethodDispatcher.dispatch(AbstractJavaResourceMethodDispatcher.java:79)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.server.model.ResourceMethodInvoker.invoke(ResourceMethodInvoker.java:469)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.server.model.ResourceMethodInvoker.apply(ResourceMethodInvoker.java:391)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.server.model.ResourceMethodInvoker.apply(ResourceMethodInvoker.java:80)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.server.ServerRuntime$1.run(ServerRuntime.java:253)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.internal.Errors$1.call(Errors.java:248)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.internal.Errors$1.call(Errors.java:244)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.internal.Errors.process(Errors.java:292)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.internal.Errors.process(Errors.java:274)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.internal.Errors.process(Errors.java:244)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.process.internal.RequestScope.runInScope(RequestScope.java:265)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.server.ServerRuntime.process(ServerRuntime.java:232)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.server.ApplicationHandler.handle(ApplicationHandler.java:679)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.servlet.WebComponent.serviceImpl(WebComponent.java:392)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.servlet.WebComponent.service(WebComponent.java:346)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.servlet.ServletContainer.service(ServletContainer.java:365)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.servlet.ServletContainer.service(ServletContainer.java:318)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.glassfish.jersey.servlet.ServletContainer.service(ServletContainer.java:205)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.servlet.ServletHolder.handle(ServletHolder.java:852)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.servlet.ServletHandler.doHandle(ServletHandler.java:544)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.server.handler.ScopedHandler.nextHandle(ScopedHandler.java:233)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.server.session.SessionHandler.doHandle(SessionHandler.java:1581)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.server.handler.ScopedHandler.nextHandle(ScopedHandler.java:233)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.server.handler.ContextHandler.doHandle(ContextHandler.java:1307)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.server.handler.ScopedHandler.nextScope(ScopedHandler.java:188)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.servlet.ServletHandler.doScope(ServletHandler.java:482)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.server.session.SessionHandler.doScope(SessionHandler.java:1549)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.server.handler.ScopedHandler.nextScope(ScopedHandler.java:186)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.server.handler.ContextHandler.doScope(ContextHandler.java:1204)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.server.handler.ScopedHandler.handle(ScopedHandler.java:141)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.server.handler.ContextHandlerCollection.handle(ContextHandlerCollection.java:221)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.server.handler.StatisticsHandler.handle(StatisticsHandler.java:173)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.server.handler.HandlerWrapper.handle(HandlerWrapper.java:127)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.server.Server.handle(Server.java:494)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.server.HttpChannel.handle(HttpChannel.java:374)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.server.HttpConnection.onFillable(HttpConnection.java:268)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.io.AbstractConnection$ReadCallback.succeeded(AbstractConnection.java:311)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.io.FillInterest.fillable(FillInterest.java:103)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.io.ChannelEndPoint$2.run(ChannelEndPoint.java:117)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.util.thread.strategy.EatWhatYouKill.runTask(EatWhatYouKill.java:336)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.util.thread.strategy.EatWhatYouKill.doProduce(EatWhatYouKill.java:313)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.util.thread.strategy.EatWhatYouKill.tryProduce(EatWhatYouKill.java:171)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.util.thread.strategy.EatWhatYouKill.run(EatWhatYouKill.java:129)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.util.thread.ReservedThreadExecutor$ReservedThread.run(ReservedThreadExecutor.java:367)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.util.thread.QueuedThreadPool.runJob(QueuedThreadPool.java:782)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.util.thread.QueuedThreadPool$Runner.run(QueuedThreadPool.java:918)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at java.base/java.lang.Thread.run(Thread.java:834)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m Caused by: java.net.SocketTimeoutException: Connect Timeout
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at org.eclipse.jetty.io.ManagedSelector$Connect.run(ManagedSelector.java:802)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:304)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
Apr 24 14:11:14 node01 docker-compose[57276]: #033[36mconnect_1  |#033[0m #011... 1 more

编辑:从最初发布此问题以来我所学到的知识,很明显这不是错误或未记录的行为,而是与 docker 网络以及容器和单独主机之间的一般网络有关。但是,我仍然不清楚如何正确配置它,即使是一次性的。我们正在使用 nginx,我们在 2 个节点前面有一个 F5 负载均衡器。我可以从任一容器 ping 另一台主机本身,因此主机至少可以相互通信。

【问题讨论】:

  • 您能否编辑您的问题以包含您的两个工人的工人配置

标签: apache-kafka apache-kafka-connect debezium


【解决方案1】:

已记录...只有一位领导者将请求转发到底层配置/状态主题。类似于任何副本只有一个领导主题分区。

除非您为它找到一个开放的 JIRA 或尝试配置与您的问题相关的每个属性,否则没有什么是错误。

特别是,您似乎没有设置 rest.advertised.listener (or the advertised host names) 以允许每台服务器相互广播自己

【讨论】:

  • 我提出这个问题的全部原因是,我相信除了 docker 和 kafka connect 方面的专家之外,文档对任何人来说都很少见。特别是,我几乎找不到任何关于在不同主机上运行的容器化 kafka 连接集群的讨论。但是您的回答与rmoff.net/2019/11/22/… 一起有助于最终获得我的问题的完整答案。所以谢谢你。
  • Docker 只是增加了一个网络层。需要正确配置的是您的网络(请注意,Kafka Operators 和 Helm Charts 会为您执行此操作),并且 Apache Kafka 项目与 Docker 没有重叠,因此“稀疏”可能不是我会使用的词。就像我说的,Kubernetes(通过 Operators 或 Helm)可能是最接近的“公共”示例,您会发现在不同主机上运行容器(尽管我在内部在 Mesos 和 Nomad 上这样做)。跨度>
  • 我同意你的评价。我通过阅读看到确实像 Kubernetes 这样的东西会是一个解决方案。但我也对如何在没有完整编排平台的情况下让其他工作人员访问每个 rest.advertised.host.name 感兴趣。我想做的是按摩我的问题,一旦说完,就让它成为一个更清晰、更有帮助的问题/答案。目前比如一个是171.15.0.2,另一个是171.28.0.2,不能互相访问。
  • 不要使用 IP。使用服务名称链接docs.docker.com/compose/networking
  • 我认为我不能使用链接,因为容器位于不同的主机上。我想我需要创建一个集群:docs.docker.com/engine/swarm。这就是为什么我正在考虑绕过 docker 网络并使用 IP 的原因——对于我们部署 swarm 来说,这是一个重大的重新架构。当前的解决方案只是暂时的,直到我们可以将这个容器集群移动到 kubernetes,我们还没有完全准备好。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2022-06-19
  • 2018-03-26
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-08-15
  • 2014-04-11
相关资源
最近更新 更多