You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

JMS应用无法消费Oracle分片队列消息的问题排查求助

Oracle分片AQ队列无法被Java应用消费的问题排查

背景

系统使用普通非分片Oracle Advanced Queue(AQ)可正常收发消息,但切换到Sharded Queue后,Java应用无法消费PL/SQL发送的消息。数据库队列表中可见消息,但Spring JMS的MessageConsumer.receive调用无法获取消息。

正常工作配置

创建队列

begin
  dbms_aqadm.create_queue_table(
    queue_table => 'poc_queue_source.nps_transactions_qt',
    queue_payload_type => 'SYS.AQ$_JMS_TEXT_MESSAGE',
    multiple_consumers => FALSE
  );
  dbms_aqadm.create_queue(
    queue_name => 'poc_queue_source.nps_transactions_queue',
    queue_table => 'poc_queue_source.nps_transactions_qt'
  );
  dbms_aqadm.start_queue(
    queue_name => 'poc_queue_source.nps_transactions_queue'
  );
end;
/

PL/SQL发送消息(触发器中)

declare
  row_json            varchar2(4000);
  enqueue_options     dbms_aq.enqueue_options_t;
  message_properties  dbms_aq.message_properties_t;
  message_handle      RAW(16);
  message             SYS.AQ$_JMS_TEXT_MESSAGE := SYS.AQ$_JMS_TEXT_MESSAGE.construct;
begin
  row_json := json_object(
    'transactionId'  value  :NEW.transactionid,
    'srcaccountid'  value  :NEW.srcaccountid,
    'dstaccountid'  value  :NEW.dstaccountid
    -- 其他字段
    null on null
  );
  message.set_text(row_json);
  dbms_aq.enqueue(
    queue_name => 'poc_queue_source.nps_transactions_queue',
    enqueue_options => enqueue_options,
    message_properties => message_properties,
    payload => message,
    msgid => message_handle
  );
end;

Spring Boot JMS接收配置

接收端为Spring Boot应用,使用Spring JMS,JMSConfiguration类相关配置如下:

@Bean
public DefaultMessageListenerContainer jmsListenerContainer(
    DefaultJmsListenerContainerFactory jmsListenerContainerFactory,
    JmsMessageListenerImpl messageListener,
    @Value("${application.jms.listener.queue-name}") String queueName,
    @Value("${application.jms.listener.idle-receives-per-task-limit}") Integer idleReceivesLimit
) {
    SimpleJmsListenerEndpoint endpoint = new SimpleJmsListenerEndpoint();
    endpoint.setMessageListener(messageListener);
    endpoint.setDestination(queueName);

    DefaultMessageListenerContainer container = jmsListenerContainerFactory.createListenerContainer(endpoint);
    // 队列空闲时缩减并发数
    container.setIdleReceivesPerTaskLimit(idleReceivesLimit);
    return container;
}

其中JmsMessageListenerImpl是自定义实现JMS MessageListener接口的组件。

依赖库

<dependency>
    <groupId>com.oracle.database.messaging</groupId>
    <artifactId>aqapi-jakarta</artifactId>
    <version>23.2.1.0</version>
</dependency>

此配置下,消息可正常收发。

失败配置(分片队列)

分片队列创建

begin
  dbms_aqadm.create_sharded_queue(
    queue_name => 'poc_queue_source.nps_transactions_shdqueue',
    multiple_consumers => FALSE,
    queue_payload_type => DBMS_AQADM.JMS_TYPE
  );
  dbms_aqadm.start_queue(
    queue_name => 'poc_queue_source.nps_transactions_shdqueue'
  );
end;
/

PL/SQL发送消息

仅修改队列名称,其余逻辑与正常配置一致:

declare
  row_json            varchar2(4000);
  enqueue_options     dbms_aq.enqueue_options_t;
  message_properties  dbms_aq.message_properties_t;
  message_handle      RAW(16);
  message             SYS.AQ$_JMS_TEXT_MESSAGE := SYS.AQ$_JMS_TEXT_MESSAGE.construct;
begin
  row_json := json_object(
    'transactionId'  value  :NEW.transactionid,
    'srcaccountid'  value  :NEW.srcaccountid,
    'dstaccountid'  value  :NEW.dstaccountid
    -- 其他字段
    null on null
  );
  message.set_text(row_json);
  dbms_aq.enqueue(
    queue_name => 'poc_queue_source.nps_transactions_queue_shdqueue',
    enqueue_options => enqueue_options,
    message_properties => message_properties,
    payload => message,
    msgid => message_handle
  );
end;

同时已更新JMS应用中Destination的队列名称。

现象

数据库队列表中存在消息,但Spring JMS的org.springframework.jms.support.destination.JmsDestinationAccessor#receiveFromConsumer方法内的MessageConsumer.receive调用无法获取消息。


问题排查与修复建议

1. 校验队列负载类型匹配

非分片队列使用SYS.AQ$_JMS_TEXT_MESSAGE作为负载类型,而分片队列创建时指定的是通用的DBMS_AQADM.JMS_TYPE,可能存在类型兼容问题。建议显式指定分片队列的负载类型与非分片队列一致:

begin
  dbms_aqadm.create_sharded_queue(
    queue_name => 'poc_queue_source.nps_transactions_shdqueue',
    multiple_consumers => FALSE,
    queue_payload_type => 'SYS.AQ$_JMS_TEXT_MESSAGE'
  );
  dbms_aqadm.start_queue(
    queue_name => 'poc_queue_source.nps_transactions_shdqueue'
  );
end;
/

2. 修正PL/SQL中的队列名称

分片队列创建名称为poc_queue_source.nps_transactions_shdqueue,但PL/SQL发送代码中写的是poc_queue_source.nps_transactions_queue_shdqueue(多了queue前缀),请确保队列名称完全一致:

dbms_aq.enqueue(
    queue_name => 'poc_queue_source.nps_transactions_shdqueue', -- 与创建的队列名称匹配
    ...
);

3. 调整JMS连接工厂配置

Oracle分片AQ队列要求客户端使用Oracle专属的JMS连接工厂,并可能需要指定分片键或启用分片支持。示例配置如下:

@Bean
public ConnectionFactory connectionFactory() throws Exception {
    AQjmsFactory factory = AQjmsFactory.getOracleAQConnectionFactory();
    Properties props = new Properties();
    props.put("user", "你的数据库用户名");
    props.put("password", "你的数据库密码");
    props.put("url", "jdbc:oracle:thin:@你的分片数据库地址");
    // 若需指定分片键,添加以下配置
    // props.put("oracle.jms.shardkey", "目标分片键");
    return factory.createConnectionFactory(props);
}

4. 确认权限配置

分片队列的权限模型与普通队列不同,需确保JMS客户端用户拥有足够权限:

-- 授予队列的收发权限
GRANT ENQUEUE, DEQUEUE ON poc_queue_source.nps_transactions_shdqueue TO 你的JMS用户;
-- 授予队列元数据查询权限
GRANT SELECT ON SYS.DBA_QUEUES TO 你的JMS用户;

5. 调整Spring JMS监听容器配置

普通队列的监听容器配置可能不适用于分片队列,建议尝试简化配置测试:

  • 暂时替换DefaultMessageListenerContainer为SimpleMessageListenerContainer,排除并发配置的影响
  • 确保监听容器使用的是Oracle专属的JMS连接工厂,而非通用JMS实现

内容的提问来源于stack exchange,提问作者thomjah

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.15 02:22:33