【问题标题】:Can you programmatically alter a queue's "dead letter" handling in a Java embedded broker?您可以在 Java 嵌入式代理中以编程方式更改队列的“死信”处理吗?
【发布时间】:2014-09-25 22:45:33
【问题描述】:

背景

概括地说,我有一个 Java 应用程序,其中某些事件应该触发为当前用户执行的某个操作。但是,事件可能非常频繁,并且动作总是相同的。因此,当第一个事件发生时,我想将动作安排在不久的将来某个时间点(例如 5 分钟)。在该时间窗口内,后续事件不应执行任何操作,因为应用程序会看到已经安排了一个操作。一旦计划的操作执行完毕,我们就返回到第 1 步,下一个事件将重新开始循环。

我的想法是通过在应用程序本身中嵌入一个内存中的 ActiveMQ 实例来实现这种过滤和节流机制(我不关心队列持久性)。

我相信 JMS 2.0 支持延迟传递的概念,延迟消息位于“暂存队列”中,直到需要传递到真正的目的地。但是,我也相信 ActiveMQ 还不支持 JMS 2.0 规范......所以我正在考虑使用生存时间 (TTL) 值和死信队列 (DLQ) 处理来模仿相同的行为。

基本上,我的消息生产者代码会将消息放在一个虚拟暂存队列中任何消费者都不会从中提取任何东西。消息将被放置一个 5 分钟的 TTL 值,并且在到期时 ActiveMQ 会将它们转储到 DLQ 中。 那是我的消息消费者实际从中消费消息的队列。

问题

我不认为我想真正从“默认”DLQ 中消费,因为我不知道 ActiveMQ 可能会在那里转储哪些与我的应用程序代码完全无关的内部事物。所以我认为最好让我的虚拟暂存队列拥有自己的自定义 DLQ。我只见过one page of ActiveMQ documentation which discusses DLQ config,它只处理独立 ActiveMQ 安装的 XML 配置文件(不是嵌入在应用程序中的内存中代理)。

是否可以在运行时为嵌入式 ActiveMQ 实例中的队列以编程方式配置自定义 DLQ?

如果您认为我走错了路,我也很想听听其他建议。我对 JMS 比对 AMQP 更熟悉,所以我不知道使用 Qpid 或其他一些可嵌入 Java 的 AMQP 代理是否更容易。无论 Apache Camel 实际上是什么(!),我相信它应该在这类事情上表现出色,但是对于这个用例来说,学习曲线可能是过大的。

【问题讨论】:

  • +1。优秀的问题写得很好。并不是每天都会看到这样的问题。

标签: java multithreading jms apache-camel activemq


【解决方案1】:

虽然您担心 Camel 对于这个用例来说可能是大材小用,但我认为 ActiveMQ 对于您所描述的用例来说已经大材小用了。

您希望安排某件事在事件发生 5 分钟后发生,并让它只消耗第一个事件并忽略第一个事件和 5 分钟结束之间的所有事件,对吗?为什么不通过ScheduledExecutorService 或您最喜欢的调度机制将您的处理方法从现在开始安排5 分钟,并将事件保存在HashMap<User, Event> 成员变量中。如果在处理方法触发之前该用户有更多事件发生,您只会看到您已经存储了一个事件并且没有存储新事件,因此您将忽略除第一个之外的所有事件。在您的处理方法结束时,从您的HashMap 中删除该用户的事件,然后将存储和安排下一个进入的事件。

仅仅为了获得这种行为而运行 ActiveMQ 似乎远远超出了您的需要。如果不是,你能解释一下原因吗?

编辑:

如果您确实走这条路,请不要使用消息 TTL 来使您的消息过期;只需让(一个也是唯一的)消费者将它们读入内存并使用上述内存解决方案每 5 分钟仅处理(最多)一批。要么有一个带有消息选择器的队列,要么使用动态队列,每个用户一个。您不需要 DLQ 来实现延迟,即使您可以让它这样做,它也不会给您批量处理所有内容的功能,因此您每 5 分钟只运行一次。这不是你想走的路,即使你知道怎么走。

【讨论】:

  • 是的!我的第一次尝试是努力使用ScheduledThreadPoolExecutor,正如您所描述的那样。问题是我还需要一种方法让后续事件检查待处理队列,并且如果该用户已经有待处理的任务,则不要为给定用户添加其他任务。在ScheduledThreadPoolExecutor 将任务放入其队列之前,它将它们包装在私有嵌套类包装器中。这种嵌套的包装器类型在ScheduledThreadPoolExecutor... 之外甚至不可见,即使是这样,它也不会暴露我最初传递的包装对象。 (1 / 3)
  • 因此,虽然嵌入 ActiveMQ 可能看起来有点矫枉过正,但我​​的替代方案是编写我自己的自定义替代方案来替代 ScheduledThreadPoolExector... 我相信这更加矫枉过正!我花了几个小时尝试这个。 (3 个中的 2 个)
  • 另外,使用真正的 MQ 代理给了我未来的发展空间。如果这个应用程序达到我想在集群中运行它的多个实例的地步,那么我可以将 ActiveMQ 代理从嵌入式方法移动到单独的独立安装......对我的应用程序代码的更改最少。
  • 更新:哦,我错过了关于“并将事件保存在成员变量中”的部分。这实际上可能是一个可行的解决方案,前提是我有一种机制可以在从池中提取任务时排出该成员变量。我会考虑一下。 (谢谢!)但是,出于我之前描述的未来可扩展性的目的,我仍然对如何使用消息代理进行操作感到好奇。
  • 顺便说一句,如果您想要一种稍微不那么脆弱的方式来处理当前安排的处理事件,您可以将 Future 而不是事件对象存储在 HashMap 中,并且然后检查Future 的状态。如果它完成了,你知道你需要安排和存储另一个。您仍然需要适当的同步,但您不必在完成处理后手动从 HashMap 中删除事件。
【解决方案2】:

一个简单的解决方案是跟踪并发结构中的待处理操作并使用ScheduledExecutorService 来执行它们:

private static final Object RUNNING = new Object();
private final ConcurrentMap<UserId, Object> pendingActions = 
    new ConcurrentHashMap<>();
private ScheduledExecutorService ses = Executors.newScheduledThreadPool(10);


public void takeAction(final UserId id) {
    Object running = pendingActions.putIfAbsent(id, RUNNING);  // atomic
    if(running == null) {                // no pending action for this user
        ses.schedule(new Runnable() {
            @Override
            public void run() {
                doWork();
                pendingActions.remove(id);
            }
        }, 5, TimeUnit.MINUTES);
    }
}

【讨论】:

    【解决方案3】:

    使用 Camel,这可以通过带有参数 completionIntervalAggregator 组件轻松实现,因此每五分钟您可以检查列表聚合消息是否为空,如果它没有向负责的路由发送消息您的用户操作并清空列表。您确实需要维护整个交换列表,只需要状态(计划或未计划用户操作)。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-04-06
      • 2014-11-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-04-01
      • 1970-01-01
      相关资源
      最近更新 更多