如何在Spring中集成Oracle AQ的JMS接口?

来源:站长查询作者:新加坡程序员头衔:程序员
导读:本期聚焦于新加坡程序员创作的《如何在Spring中集成Oracle AQ的JMS接口?》,敬请观看详情。如果项目消息中间件选型落在Oracle数据库,很可能会用到Oracle AQ。把AQ的JMS接口接入Spring时,核心不是简单引入依赖,而是要让连接工厂、目的地对象和监听容器三者正确协作。本文围绕Spring整合Oracle AQ JMS的完整链路展开,先说明aqapi与ojdbc等必要依赖,再演示通过AQjmsFactory获取QueueConnectionFactory以及使用JmsTemplate完成点对点收发。异步消费部分会重点分析DefaultMessageListenerContainer的配置方式,包括并发数、会话事务和错误回滚。最后讨论Oracle AQ在本地事务中的特殊表现,以及如何避免消息重复消费和确认丢失。读完可以快速搭建一个稳定可用的Spring与Oracle AQ消息收发工程。

Oracle AQ(Advanced Queuing)是Oracle数据库内置的消息队列组件,它把消息直接存储在数据库表中,并提供了完整的JMS 1.1接口。对很多以Oracle为核心存储的Java项目来说,不需要额外部署RabbitMQ或Kafka,就能通过AQ实现异步解耦、任务缓冲和事件通知。Spring的spring-jms模块对JMS封装得比较完善,但Oracle AQ的连接工厂和目的地解析与常见ActiveMQ、IBM MQ略有差异,集成时容易在配置环节出错。

如何在Spring中集成Oracle AQ的JMS接口?

依赖与数据库对象准备

要在Spring中使用Oracle AQ的JMS接口,第一步是引入Oracle官方提供的JMS API实现。核心依赖是aqapi,它包含oracle.jms.AQjmsFactory等类;还需要ojdbc数据库驱动,因为AQ的底层连接完全依赖JDBC会话。spring-jms提供JmsTemplate和DefaultMessageListenerContainer等模板和容器支持。典型Maven配置如下:

<dependency>
    <groupId>com.oracle.database.messaging</groupId>
    <artifactId>aqapi</artifactId>
    <version>21.3.0.0</version>
</dependency>
<dependency>
    <groupId>com.oracle.database.jdbc</groupId>
    <artifactId>ojdbc8</artifactId>
    <version>21.3.0.0</version>
</dependency>
<dependency>
    <groupId>org.springframework</groupId>
    <artifactId>spring-jms</artifactId>
</dependency>

依赖就绪后,还需要在Oracle数据库中创建队列表和队列。Oracle AQ的消息并不是简单的内存对象,而是持久化在数据库表里。对于JMS文本消息,可以使用SYS.AQ$_JMS_TEXT_MESSAGE作为负载类型。下面的PL/SQL块创建了一张队列表、一个队列并启动队列:

BEGIN
  DBMS_AQADM.CREATE_QUEUE_TABLE(
    queue_table => 'MY_QUEUE_TABLE',
    queue_payload_type => 'SYS.AQ$_JMS_TEXT_MESSAGE',
    multiple_consumers => FALSE
  );
  DBMS_AQADM.CREATE_QUEUE(
    queue_name => 'MY_QUEUE',
    queue_table => 'MY_QUEUE_TABLE'
  );
  DBMS_AQADM.START_QUEUE(queue_name => 'MY_QUEUE');
END;
/

这里创建的是单消费者队列,也就是点对点模式。如果业务需要多个消费者同时订阅并各自收到一份消息,需要把multiple_consumers设置为TRUE,并使用Topic或者JMS的发布订阅模型,但多数Spring集成场景先跑通Queue类型更简单。

配置QueueConnectionFactory与JmsTemplate

Oracle AQ没有独立的JNDI服务,它的连接工厂不能像ActiveMQ那样从Broker URL创建。正确做法是通过oracle.jms.AQjmsFactory的静态工厂方法,基于已有的DataSource创建QueueConnectionFactory。Spring容器中可以直接暴露一个@Bean,把DataSource当作入参。

@Configuration
public class AqJmsConfig {

    @Bean
    public QueueConnectionFactory queueConnectionFactory(DataSource dataSource) throws JMSException {
        return AQjmsFactory.getQueueConnectionFactory(dataSource);
    }

    @Bean
    public Queue myQueue(DataSource dataSource) throws JMSException {
        return AQjmsFactory.getQueue(dataSource, "MY_QUEUE");
    }

    @Bean
    public JmsTemplate jmsTemplate(QueueConnectionFactory queueConnectionFactory, Queue myQueue) {
        JmsTemplate jmsTemplate = new JmsTemplate(queueConnectionFactory);
        jmsTemplate.setDefaultDestination(myQueue);
        return jmsTemplate;
    }
}

如果要沿用传统的XML配置方式,也是常见的做法。XML里同样通过工厂方法创建Bean,注意factory-method和constructor-arg的配合。配置示例如下:

<bean id="queueConnectionFactory" class="oracle.jms.AQjmsFactory" factory-method="getQueueConnectionFactory">
    <constructor-arg ref="dataSource"/>
</bean>
<bean id="myQueue" class="oracle.jms.AQjmsFactory" factory-method="getQueue">
    <constructor-arg ref="dataSource"/>
    <constructor-arg value="MY_QUEUE"/>
</bean>
<bean id="jmsTemplate" class="org.springframework.jms.core.JmsTemplate">
    <property name="connectionFactory" ref="queueConnectionFactory"/>
    <property name="defaultDestination" ref="myQueue"/>
</bean>

这里有一个容易忽略的点:AQjmsFactory.getQueueConnectionFactory返回的是javax.jms.QueueConnectionFactory,它本身继承自javax.jms.ConnectionFactory,因此可以直接赋给JmsTemplate。getQueue方法同样需要DataSource和队列名,返回的Queue对象可以直接作为JmsTemplate的默认目的地,避免了目的地名称查找失败的问题。

发送消息与同步接收

JmsTemplate封装了连接创建、会话打开和资源关闭,发送文本消息最直接的方式是使用send方法配合lambda表达式。lambda的MessageCreator接口负责创建一条文本消息,代码很简洁。

@Service
public class MessageProducer {

    private final JmsTemplate jmsTemplate;

    public MessageProducer(JmsTemplate jmsTemplate) {
        this.jmsTemplate = jmsTemplate;
    }

    public void sendText(String payload) {
        jmsTemplate.send(session -> session.createTextMessage(payload));
    }
}

如果发送的是对象消息,可以使用convertAndSend方法,由Spring自动选择MessageConverter。但Oracle AQ的JMS实现中,对象消息底层使用Oracle的ADT类型或序列化对象,跨应用版本时可能出现反序列化兼容性问题,所以生产环境更推荐发送TextMessage或BytesMessage,再由业务自行序列化为JSON。

同步接收可以通过receive方法阻塞等待一条消息,设置receiveTimeout可以避免无限等待。下面的代码展示了如何读取文本消息并返回字符串:

public String receiveText() throws JMSException {
    Message message = jmsTemplate.receive();
    if (message instanceof TextMessage) {
        TextMessage textMessage = (TextMessage) message;
        return textMessage.getText();
    }
    return null;
}

同步接收适合后台任务或管理接口临时拉取消息,对于需要持续消费的场景则不建议使用while循环轮询,资源和性能都不如异步监听容器。Oracle AQ的消息存储在数据库表中,频繁同步接收会造成额外的SQL查询压力。

配置异步监听容器

Spring JMS的异步消费核心是DefaultMessageListenerContainer,它会维护一组JMS Session并注册MessageListener,收到消息后回调onMessage方法。对于Oracle AQ,连接工厂和目的地都应使用前面定义的Bean,不建议使用destinationName加JNDI查找。

@Bean
public DefaultMessageListenerContainer messageListenerContainer(
        QueueConnectionFactory queueConnectionFactory,
        Queue myQueue,
        MessageListener aqMessageListener) {

    DefaultMessageListenerContainer container = new DefaultMessageListenerContainer();
    container.setConnectionFactory(queueConnectionFactory);
    container.setDestination(myQueue);
    container.setMessageListener(aqMessageListener);
    container.setSessionTransacted(true);
    container.setConcurrentConsumers(1);
    container.setMaxConcurrentConsumers(5);
    container.setReceiveTimeout(1000L);
    return container;
}

setSessionTransacted设置为true后,监听容器会在消息处理成功后提交JMS会话,如果监听器抛出运行时异常则回滚,消息重新进入队列等待下一次投递。setConcurrentConsumers和setMaxConcurrentConsumers控制并发消费线程范围,Oracle AQ底层是数据库表,并发过高会增加锁竞争,一般从1到5开始压测再逐步调整。setReceiveTimeout表示在没有消息时等待数据库返回的毫秒数,适当调大能减少空轮询。

监听器的实现需要实现javax.jms.MessageListener接口。处理逻辑中要捕获JMSException并转换为运行时异常,否则容器无法感知失败,消息可能被错误确认。下面是一个简单的文本消息监听器:

@Component
public class AqMessageListener implements MessageListener {

    @Override
    public void onMessage(Message message) {
        if (message instanceof TextMessage) {
            try {
                String text = ((TextMessage) message).getText();
                process(text);
            } catch (JMSException e) {
                throw new RuntimeException("处理Oracle AQ消息失败", e);
            }
        }
    }

    private void process(String text) {
        // 业务处理逻辑
        System.out.println(text);
    }
}

如果业务处理失败且希望消息重试,可以抛出运行时异常;但要注意Oracle AQ默认不会限制重试次数,如果消息本身有数据问题,会一直重复消费。可以在监听器里增加重试计数或使用后文提到的异常表机制。

事务一致性与调优建议

Oracle AQ与普通JMS中间件最大的区别是它的消息和业务表同在一个数据库实例中。这个特性有好处,也有麻烦。好处是队列操作天然可以纳入数据库事务,麻烦是JMS本地事务和Spring的JDBC事务并不是自动合并的。很多项目遇到的问题是:业务表写入成功但消息发送失败,或者消息消费成功但业务表回滚。

对于发送端,如果业务更新和消息入队需要强一致,建议把这两步放在同一个数据库本地事务中执行。可以在Service方法上使用DataSourceTransactionManager管理JDBC事务,并确保JmsTemplate使用的DataSource连接与当前事务绑定。由于Oracle AQ的JMS会话底层依赖同一JDBC连接,这种方式在正确配置后可以做到同时提交或回滚。简化方案则是只给JMS设置本地事务,业务表更新失败时通过补偿任务处理,避免引入复杂的XA配置。

消费端建议保持sessionTransacted为true,这样消息确认、业务处理和数据库写操作共享一个JMS本地事务。监听器里如果调用了其他Service方法写数据库,异常抛出后消息会自动回滚。但要注意不能手动调用message.acknowledge,否则会破坏容器的事务管理。还可以为容器设置errorHandler,把反复失败的异常记录到日志或异常表,方便人工排查。

调优方面,Oracle AQ的监听容器不适合设置过大的并发数,因为消息存储在数据库表中,每个消费者都会执行SELECT和UPDATE操作,过高并发容易造成enqueue/dequeue相关的锁等待。可以从1到3个并发开始压测,观察GV$AQ视图和AWR报告中的等待事件。另外setCacheLevel可以控制JMS Session的缓存级别,默认值在大多数场景够用,如果发现频繁创建数据库Session再考虑调整。

Oracle AQJMSSpring集成修改时间:2026-10-06 04:00:56

免责声明:已尽一切努力确保本网站所含信息的准确性。网站作品多为原创整理与精心创作,观点力求客观中立。本站旨在免费分享,内容仅供个人学习、研究或参考使用。若引用了第三方作品,版权归原作者所有。如内容涉及您的权益,请联系我们进行处理Email:chomcom@qq.com。
引用或转载本作品时,请注明当前出处:https://www.ipipp.com/html/1006/66286.html,基于非商业用途的前提下,欢迎转载或二创本作品。
内容垂直聚焦
专注技术核心技术栏目,确保每篇文章深度聚焦于实用技能。从代码技巧到架构设计,为用户提供无干扰的纯技术知识沉淀,精准满足专业提升需求。
知识结构清晰
覆盖从开发到部署的全链路。AI、前端、编程、数据库、服务器、建站、系统层层递进,构建清晰学习路径,帮助用户系统化掌握开发与运维所需的核心技术。
深度技术解析
拒绝泛泛而谈,深入技术细节与实践难点。无论是数据库优化还是服务器配置,均结合真实场景与代码示例进行剖析,致力于提供可直接应用于工作的解决方案。
专业领域覆盖
精准对应开发生命周期。从前端界面到后端编程,从数据库操作到服务器运维,形成完整闭环,一站式满足全栈工程师和运维人员的技术需求。
即学即用高效
内容强调实操性,步骤清晰、代码完整。用户可根据教程直接复现和应用于自身项目,显著缩短从学习到实践的距离,快速解决开发中的具体问题。
持续更新保障
专注既定技术方向进行长期、稳定的内容输出。确保各栏目技术文章持续更新迭代,紧跟主流技术发展趋势,为用户提供经久不衰的学习价值。