【问题标题】:Receiving multiple messages from MQ asynchronously从 MQ 异步接收多条消息
【发布时间】:2019-08-02 02:02:32
【问题描述】:

我在我的应用程序中使用 Spring + Hibernate + JPA。

我需要从Websphere MQ 读取消息并将消息插入数据库。 有时可能会有连续的消息可用,有时消息的数量会非常少,有时我们可能无法期待来自Queue 的消息。

目前我正在一一阅读消息并将它们插入数据库。但在性能方面并没有多大帮助。

我的意思是当我有大量消息时(例如队列中的 30 万条消息),我无法更快地插入它们。每秒插入数据库的实体数量不是很高。因为我确实为每一个实体做出了承诺。

我想使用休眠批处理,以便我可以在单个提交中插入实体列表。 (例如:每次提交 30 到 40 条消息)

问题:

  1. 如何从队列接收多条消息? (我已经检查过 BatchMessageListenerContainer 可能会有所帮助。但我无法获得一些参考)

  2. 我应该将 db 插入过程与 onMessage 方法分开吗?那么该线程将被释放到池中并可用于从队列中挑选下一条消息?

  3. 并行线程使用情况?

当前实现:

消息监听器:

<bean id="myMessageListener" class="org.mypackage.MyMessageListener">

<bean id="jmsContainer" class="org.springframework.jms.listener.DefaultMessageListenerContainer">
    <property name="connectionFactory" ref="connectionFactory"/>
    <property name="destinationName" ref="queue"/>
    <property name="messageListener" ref="myMessageListener"/>
    <property name ="concurrentConsumers" value ="10"/>
    <property name ="maxConcurrentConsumers" value ="50"/>        
</bean>

监听类:

package org.mypackage.MyMessageListener;

import javax.jms.JMSException;
import javax.jms.Message;
import javax.jms.MessageListener;

import org.mypackage.service.MyService;

public class MyMessageListener implements MessageListener {

    @Autowired
    private MyService myService;

    @Override
    public void onMessage(Message message) {
        try {
             TextMessage textMessage = (TextMessage) message;
             // parse the message
             // Process the message to DB
        } catch (JMSException e1) {
             e1.printStackTrace();
        }
    }
}

【问题讨论】:

  • 为了能够批量插入,您需要一个要在一批中插入的项目列表。我认为您必须扩展队列以根据时间和/或大小为您交付一批订单。
  • @FlorianDe 在我的情况下,队列的发件人无法根据时间和/或大小发送一批订单。
  • 我认为您必须根据mcve 更新您的问题。因为目前看不到您的实施细节时,很难阐明任何解决方案。
  • @FlorianDe,我通过添加我的实现代码更新了问题

标签: multithreading parallel-processing spring-batch spring-jms mq


【解决方案1】:

不清楚您的要求是什么。

Spring Batch 项目提供了一个BatchMessageListenerContainer

消息侦听器容器适用于通过配置提供的建议来拦截消息接收。 要在单个事务中启用消息批处理,请在建议链中使用 TransactionInterceptor 和 RepeatOperationsInterceptor(在基类中设置或不设置事务管理器)。然后容器将使用RepeatOperations 在同一个线程中接收多条消息,而不是接收一条消息并对其进行处理。与 RepeatOperations 和事务拦截器一起使用。如果事务拦截器使用 XA,则使用 XA 连接工厂或 TransactionAwareConnectionFactoryProxy 将 JMS 会话与正在进行的事务同步(打开失败后重复消息的可能性)。在后一种情况下,您不需要在基类中提供事务管理器 - 它只会在途中阻止 JMS 会话与数据库事务同步。

【讨论】:

  • 我已经更新了我的问题。我正在使用 Websphere MQ,我想接收大量消息,以便我可以在一次提交到 DB 中插入多条记录。
  • 对不起;我没有阅读标签;我回答了 RabbitMQ。请参阅我的答案的编辑。
  • 感谢您的回答。你的意思是我可以使用 BatchMessageListenerContainer 代替 DMLC 吗?我应该将 onMessage 方法参数更改为 List 吗?我无法获得 BatchMessageListenerContainer 用法的示例
  • 不,您仍然会一次收到一条消息,因此您需要累积它们;然后,当你有一个完整的批量更新数据库和事务将提交(数据库和批处理)。
  • 如果 batchSize 为 30 并且如果队列中只有 10 条消息可用,并且剩余 20 条消息可能会在一个小时后到达队列。等待剩余消息到达直到达到批量大小对我来说并不好。因为我想等待几毫秒(例如:3000ms),如果在那段时间内收到 30 条消息,那么我可以将它们更新到 DB,否则我可以将可用的消息数更新到 DB。请问这可以通过这个容器实现吗?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-04-08
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-02-02
  • 2021-09-24
相关资源
最近更新 更多