【问题标题】:kafka client is sending request to partition where broker which went downkafka 客户端正在向代理失败的分区发送请求
【发布时间】:2019-10-04 23:01:09
【问题描述】:

我正在使用 kafka-node 模块向 kafka 发送消息。 在集群环境中,我有一个具有 3 个分区和复制因子为 3 的主题。

主题描述是 -

Topic:clusterTopic      PartitionCount:3        ReplicationFactor:3    Configs:min.insync.replicas=2,segment.bytes=1073741824
        Topic: clusterTopic     Partition: 0    Leader: 1       Replicas: 1,2,3 Isr: 1,2,3
        Topic: clusterTopic     Partition: 1    Leader: 2       Replicas: 2,3,1 Isr: 1,2,3
        Topic: clusterTopic     Partition: 2    Leader: 3       Replicas: 3,1,2 Isr: 1,2,3

生产者配置 -

        "requireAcks": 1,
        "attributes": 2,
        "partitionerType": 2,
        "retries": 2

当我发送数据时,它遵循像循环方式一样的循环(2)分区类型

当我按照以下步骤操作时

  • 获取连接到 kafka:9092,kafka:9093 的 HighLevelProducer 实例
  • 发消息
  • 手动停止 kafka-server:9092
  • 尝试使用 HighLevelProducer 发送另一条消息,然后 send() 将 触发回调错误:TimeoutError: Request timed out after 30000毫秒

我期望的是,如果一个分区不可访问(因为代理关闭),生产者应该自动将数据发送到下一个可用分区,但由于异常而我丢失了消息

异常如下-

  TimeoutError: Request timed out after 3000ms
    at new TimeoutError (\package\node_modules\kafka-node\lib\errors\TimeoutError.js:6:9)
    at Timeout.timeoutId._createTimeout [as _onTimeout] (\package\node_modules\kafka-node\lib\kafkaClient.js:980:14)
    at ontimeout (timers.js:424:11)
    at tryOnTimeout (timers.js:288:5)
    at listOnTimeout (timers.js:251:5)
    at Timer.processTimers (timers.js:211:10)
(node:56416) [DEP0079] DeprecationWarning: Custom inspection function on Objects via .inspect() is deprecated
  kafka-node:KafkaClient kafka-node-client reconnecting to kafka1:9092 +3s
  kafka-node:KafkaClient createBroker kafka1 9092 +1ms
  kafka-node:KafkaClient kafka-node-client reconnecting to kafka1:9092 +3s
  kafka-node:KafkaClient createBroker kafka1 9092 +0ms

【问题讨论】:

  • 您确定复制工作正常吗?同步副本是否显示所有可用的代理?
  • 请在您的问题中添加更多详细信息,即主题的设置min.insync.replicas,以及您的生产者中的acksretriesdelivery.timeout.ms
  • @cricket_007 是的,复制工作正常(参考主题描述)

标签: node.js apache-kafka kafka-producer-api node-kafka


【解决方案1】:

请发送引导服务器进行确认,但根据手头的信息,我相信您遇到的情况如下:

  • 您已将 min.insync.replicas 设置为 2
  • 您已将 acks 设置为 1

通过这些设置,生产者会将事件发送到领导者副本并假设消息是安全的。

如果在发送后立即失败,并且在追随者赶上之前,您将丢失消息,因为您只等待一个确认。

但是,从代理的角度来看,您指定主题可用的要求是 2 个同步副本。默认情况下,仅允许同步的副本被选为领导者。由于第一个失败会导致关注者不同步,您的主题可能会被强制下线。您可以在测试中验证这一点,假设有一些设置。

要纠正,请尝试以下操作:

  1. 如果高可用性最重要,请将 min.insync.replicas 设置为 1 并将acks 设置为 1
  2. 如果不能接受数据丢失,请将 min.insync.replicas 设置为 2 并将 acks 设置为 all

您也可以将 unclean.leader.election.enable 设置为 true 以实现高可用性,因为这将允许不同步的副本被选为领导者,但这样就有可能丢失数据。

【讨论】:

  • 嘿,我认为问题不在于 ack 和 insync 副本(我尝试了您在生产者上提出的仍然相同的错误),但它是循环的 partitionerType。因此,数据被发送到分区 1,2,然后是 3,但是如果我停止存在任何一个分区的代理,则模块会尝试将其发送到相同的分区。我希望它发现下一个可用的分区/代理并重新发送数据
  • 您是否尝试将 unclean.leader.election.enable 设置为 true?如果生产者向死亡的代理发送请求,它将失败并随后尝试另一个代理,但前提是该分区有领导者。从代理的角度查看您的指标,当这种情况发生时您是否看到任何离线分区?
  • 我尝试将 unclean.leader.election.enable 设置为 true,但我仍然收到我关闭的代理的连接错误。它重试 5 次,然后因 KafkaJSNumberOfRetriesExceeded 异常而停止。上一个错误是连接错误:将 ECONNREFUSED 连接到代理
  • 你能在你的生产者配置中设置这个并发送日志吗:logLevel: logLevel.DEBUG
  • 更新了异常问题
猜你喜欢
  • 1970-01-01
  • 2019-12-21
  • 1970-01-01
  • 2016-07-31
  • 1970-01-01
  • 2020-09-17
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多