【问题标题】:Spring Cloud Stream - how to handle downstream blocks?Spring Cloud Stream - 如何处理下游块?
【发布时间】:2021-04-20 07:51:19
【问题描述】:

在我们的 Kafka 集群计划停机期间,我们基本上遇到了以下问题 How to specify timeout for sending message to RabbitMQ using Spring Cloud Stream?(显然是使用 Kafka 而不是 RabbitMQ)。

@GaryRussell 的回答:

频道sendTimeout 仅适用于频道本身可以阻塞的情况,例如带有当前已满的有界队列的QueueChannel;调用者将阻塞,直到队列中有空间可用,或者发生超时。

在这种情况下,阻塞在通道的下游,因此 sendTimeout 无关紧要(无论如何,它是一个 DirectChannel,无论如何都不能阻塞,订阅的处理程序直接在调用线程上调用)。

您看到的实际阻塞很可能在rabbitmq客户端中的socket.write()中,它没有超时且不可中断;调用线程无法执行任何操作来“超时”写入。

我知道的唯一可能的解决方案是通过在连接工厂上调用resetConnection() 来强制关闭兔子连接。

很好地解释了为什么有问题的方法 (org.springframework.integration.channel.AbstractSubscribableChannel#doSend) 没有考虑到timeout。但是,这对我来说仍然有点奇怪。

spring-integration-kafka-3.2.1.RELEASE-sources.jar!/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandler.java:566 中,我们可以看到,如果需要sync 行为:

565    if (this.sync) {
566        Long sendTimeout = this.sendTimeoutExpression.getValue(this.evaluationContext, message, Long.class);
567        if (sendTimeout == null || sendTimeout < 0) {
568            future.get();
569        }
570        else {
571            try {
572                future.get(sendTimeout, TimeUnit.MILLISECONDS);
573            }
574            catch (TimeoutException te) {
575                throw new MessageTimeoutException(message, "Timeout waiting for response from KafkaProducer", te);
576            }
577        }
578    }

被调用,其中考虑了超时。 sendTimeoutExpression 被分配了一个默认值:

    private static final long DEFAULT_SEND_TIMEOUT = 10000;
    private Expression sendTimeoutExpression = new ValueExpression<>(DEFAULT_SEND_TIMEOUT);

然而,我们的堆栈跟踪揭示了一些不同的东西:

"pool-1-thread-3" - Thread t@108
   java.lang.Thread.State: TIMED_WAITING
    at sun.misc.Unsafe.park(Native Method)
    - parking to wait for <4ebda621> (a org.springframework.util.concurrent.SettableListenableFuture$SettableTask)
    at java.util.concurrent.locks.LockSupport.parkNanos(LockSupport.java:215)
    at java.util.concurrent.FutureTask.awaitDone(FutureTask.java:426)
    at java.util.concurrent.FutureTask.get(FutureTask.java:204)
    at org.springframework.util.concurrent.SettableListenableFuture.get(SettableListenableFuture.java:134)
*   at org.springframework.integration.kafka.outbound.KafkaProducerMessageHandler.processSendResult(KafkaProducerMessageHandler.java:572)
    at org.springframework.integration.kafka.outbound.KafkaProducerMessageHandler.handleRequestMessage(KafkaProducerMessageHandler.java:414)
    at org.springframework.integration.handler.AbstractReplyProducingMessageHandler.handleMessageInternal(AbstractReplyProducingMessageHandler.java:134)
    at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:69)
    at org.springframework.cloud.stream.binder.AbstractMessageChannelBinder$SendingHandler.handleMessageInternal(AbstractMessageChannelBinder.java:1035)
    at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:69)
    at org.springframework.integration.dispatcher.AbstractDispatcher.tryOptimizedDispatch(AbstractDispatcher.java:115)
    at org.springframework.integration.dispatcher.UnicastingDispatcher.doDispatch(UnicastingDispatcher.java:133)
    at org.springframework.integration.dispatcher.UnicastingDispatcher.dispatch(UnicastingDispatcher.java:106)
    at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:72)
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:570)

标有* 的调用对应于future.get(sendTimeout, TimeUnit.MILLISECONDS); 调用。

看到底层客户端似乎支持它(鉴于future.get() 调用支持超时这一事实),如何设置?我可以在活页夹引用中找到的唯一两个属性(参见 here)是 spring.cloud.stream.kafka.binder.healthTimeoutbatchTimeout,据我所知,这不会影响此设置。

看看KafkaProducerMessageHandler是如何在私有类org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder.ProducerConfigurationMessageHandler中构造的,bean覆盖似乎不是推荐的方式。

【问题讨论】:

    标签: java spring-integration spring-kafka spring-cloud-stream


    【解决方案1】:

    似乎没有记录,但类似于侦听器容器定制器https://docs.spring.io/spring-cloud-stream/docs/3.1.2/reference/html/spring-cloud-stream.html#_advanced_consumer_configuration,您可以添加ProducerMessageHandlerCustomizer @Bean 以在消息处理程序上设置任意属性。

    在较新版本的处理程序中,超时始终配置为至少与ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG 一样多,以避免误报(在处理程序超时后发布成功)。

    【讨论】:

    • 感谢您提供的信息 - 如果可以的话,我会在接下来的几天内尝试为文档进行 PR。关于delivery.timeout.ms 设置的提示也很有价值。由于在2.5.1 Kafka 客户端(我们正在使用的)中将其设置为120000(120 秒),并且默认的 SCS 超时设置为10000(10 秒),我想假设它是安全的我们过去可能有过假阴性。我会尝试一些事情并回到这个答案。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-10-03
    • 1970-01-01
    • 2019-06-09
    • 1970-01-01
    • 2021-04-27
    • 2017-11-24
    • 1970-01-01
    相关资源
    最近更新 更多