【发布时间】:2020-11-11 15:19:38
【问题描述】:
背景
我们有一个使用 Azure 服务总线作为消息代理的数据传输解决方案。我们正在通过x 队列从x 数据集传输数据-x 专用QueueClients 作为发送者。一些发件人以每两秒一条消息的速度发布消息,而另一些发件人每 15 分钟发布一条。
数据源端(发送者所在的位置)的应用程序运行良好,为我们提供了所需的吞吐量。
另一方面,我们有一个应用程序,每个队列有一个 QueueClient 接收器,配置如下:
-
maxConcurrentCalls=1 -
autoComplete=true(如果接收模式 =RECEIVEANDDELETE)和false(如果接收模式 =PEEKLOCK) - 我们有一些接收器,如果它们意外关闭,希望将消息保留在服务总线队列。 -
maxAutoRenewDuration=3分钟(所有队列的锁定持续时间 =30秒) - 单线程执行器服务
在这些接收者中注册的MessageHandler 执行以下操作:
public CompletableFuture<Void> onMessageAsync(final IMessage message) {
// deserialize the message body
final CustomObject customObject = (CustomObject)SerializationUtils.deserialize((byte[])message.getMessageBody().getBinaryData().get(0));
// process processDB1() and processDB2() asynchronously
final List<CompletableFuture<Boolean>> processFutures = new ArrayList<CompletableFuture<Boolean>>();
processFutures.add(processDB1(customObject)); // processDB1() returns Boolean
processFutures.add(processDB2(customObject)); // processDB2() returns Boolean
// join both the completablefutures to get the result Booleans
List<Boolean> results = CompletableFuture.allOf(processFutures.toArray(new CompletableFuture[processFutures.size()])).thenApply(future -> processFutures.stream()
.map(CompletableFuture<Boolean>::join).collect(Collectors.toList())
if (results.contains(false)) {
// dead-letter the message if results contains false
return getQueueClient().deadLetterAsync(message.getLockToken());
} else {
// complete the message otherwise
getQueueClient().completeAsync(message.getLockToken());
}
}
我们测试了以下场景:
场景 1 - 接收模式 = RECEIVEANDDELETE,消息发布率:30/ minute
预期行为
应该以恒定的吞吐量连续接收消息(不一定是发布消息的源处的吞吐量)。
实际行为
我们观察到 QueueClient 随机长时间不活动 - 从几分钟到几小时不等 - 服务总线命名空间没有传出消息(在 Metrics 图表上观察到)并且没有消耗同一时间段的日志!
场景 2 - 接收模式 = PEEKLOCK,消息发布率:30/ minute
预期行为
应该以恒定的吞吐量连续接收消息(不一定是发布消息的源处的吞吐量)。
实际行为
在应用程序运行 20-30 分钟后,我们不断看到MessageLockLostException。
我们尝试执行以下操作 -
- 我们将预取计数(从 20 * 处理速率 - 如最佳实践指南中所述)减少到最低限度(甚至在一个测试周期中达到
0),以减少数量。为客户端锁定的消息数 - 将
maxAutoRenewDuration增加到 5 分钟 - 我们的processDB1()和processDB2()在几乎 90% 的情况下不会花费超过一两秒 - 所以,我认为锁定持续时间为 30 秒和 @987654346 @ 在这里不是问题。 - 移除阻塞
CompletableFuture.get()并使处理同步。
这些调整都没有帮助我们解决问题。我们观察到COMPLETE 或RENEWMESSAGELOCK 正在抛出MessageLockLostException。
我们需要帮助来寻找以下问题的答案:
-
为什么在场景 1 中
QueueClient长时间处于不活动状态? -
我们怎么知道
MessageLockLostExceptions 被抛出了,因为锁确实已经过期了?我们怀疑锁不会过早过期,因为我们的处理会在一两秒内发生。禁用预取也没有为我们解决这个问题。
版本和服务总线详细信息
- Java -
openjdk-11-jre - Azure 服务总线命名空间层:
Standard - Java SDK 版本 -
3.4.0
【问题讨论】:
标签: java azure azureservicebus messaging