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接口,第一步是引入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再考虑调整。