【问题标题】:Periods of prolonged inactivity and frequent MessageLockLostException in QueueClientQueueClient 中长时间不活动和频繁出现 MessageLockLostException
【发布时间】: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

我们尝试执行以下操作 -

  1. 我们将预取计数(从 20 * 处理速率 - 如最佳实践指南中所述)减少到最低限度(甚至在一个测试周期中达到 0),以减少数量。为客户端锁定的消息数
  2. maxAutoRenewDuration 增加到 5 分钟 - 我们的 processDB1()processDB2() 在几乎 90% 的情况下不会花费超过一两秒 - 所以,我认为锁定持续时间为 30 秒和 @987654346 @ 在这里不是问题。
  3. 移除阻塞 CompletableFuture.get() 并使处理同步。

这些调整都没有帮助我们解决问题。我们观察到COMPLETERENEWMESSAGELOCK 正在抛出MessageLockLostException

我们需要帮助来寻找以下问题的答案:

  1. 为什么在场景 1 中 QueueClient 长时间处于不活动状态
  2. 我们怎么知道MessageLockLostExceptions 被抛出了,因为锁确实已经过期了?我们怀疑锁不会过早过期,因为我们的处理会在一两秒内发生。禁用预取也没有为我们解决这个问题。

版本和服务总线详细信息

  • Java - openjdk-11-jre
  • Azure 服务总线命名空间层:Standard
  • Java SDK 版本 - 3.4.0

【问题讨论】:

    标签: java azure azureservicebus messaging


    【解决方案1】:

    问题不在于 QueueClient 对象本身。这是我们从MessageHandlerprocessDB1(customObject)processDB2(customObject) 中触发的流程。由于这些过程没有优化,消息消耗下降并且锁 gor 过期(在 peek-lock 模式下),因为处理程序花费更多时间(相对于消息发布到队列的速率)来完成这些操作.

    优化流程后,消耗和完成(在peek-lock模式下)都很好。

    【讨论】:

      【解决方案2】:

      对于场景 1:

      如果您启用了duplicate detection history,则可能会根据以下说明的情况发生此行为:

      我已启用 30 秒。我经常用重复的消息打服务总线(我的案例消息与来自客户端的相同 messageid - 30 /每分钟)。我会看到窗口没有活动。虽然消息是在服务总线上从发送客户端收到的,但我无法在传出消息中看到它们。您可能会检查您是否再次遇到被过滤的重复消息 - 进而导致传出不活动。

      另请注意:创建队列后,您无法启用/禁用重复检测。您只能在创建队列时这样做。

      【讨论】:

      • 我们没有启用重复检测。这是我们从MessageHandler 执行的进程的问题。我们优化了这些以解决这个问题。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2016-09-04
      • 2013-02-14
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-01-28
      • 2015-06-07
      相关资源
      最近更新 更多