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

如何在WSO2 MI 4.2中从RabbitMQ入站端点提取队列名、优先级及消息ID

在WSO2 MI 4.2中从RabbitMQ消息提取元数据并更新数据库状态

一、RabbitMQ入站端点配置(开启元数据传递)

首先要确保RabbitMQ入站端点启用消息属性传递,这样队列名、消息ID、优先级等元数据才会被带入消息上下文。以下是完整配置示例:

<?xml version="1.0" encoding="UTF-8"?>
<inboundEndpoint name="RabbitMQInbound" protocol="rabbitmq" sequence="ProcessRabbitMQMessage" onError="ErrorSequence" xmlns="http://ws.apache.org/ns/synapse">
    <parameters>
        <!-- RabbitMQ服务器连接配置 -->
        <parameter name="rabbitmq.server.host.name">localhost</parameter>
        <parameter name="rabbitmq.server.port">5672</parameter>
        <parameter name="rabbitmq.server.username">guest</parameter>
        <parameter name="rabbitmq.server.password">guest</parameter>
        <!-- 队列配置 -->
        <parameter name="rabbitmq.queue.name">your-target-queue</parameter>
        <parameter name="rabbitmq.queue.durable">true</parameter>
        <!-- 关键:启用消息属性传递 -->
        <parameter name="rabbitmq.message.properties.enable">true</parameter>
        <!-- 其他可选配置 -->
        <parameter name="rabbitmq.consumer.tag">mi-rabbit-consumer</parameter>
        <parameter name="rabbitmq.prefetch.count">10</parameter>
    </parameters>
</inboundEndpoint>

二、集成序列中提取元数据

WSO2 MI会将RabbitMQ元数据存储在消息上下文属性中,可通过get-property()函数直接提取。以下是序列提取逻辑的代码片段:

<?xml version="1.0" encoding="UTF-8"?>
<sequence name="ProcessRabbitMQMessage" xmlns="http://ws.apache.org/ns/synapse">
    <!-- 提取元数据到本地属性 -->
    <property name="msgId" expression="get-property('rabbitmq.message.id')" scope="default" type="STRING"/>
    <property name="msgPriority" expression="get-property('rabbitmq.message.priority')" scope="default" type="INTEGER"/>
    <property name="queueName" expression="get-property('rabbitmq.queue.name')" scope="default" type="STRING"/>

    <!-- 可选:打印元数据用于调试 -->
    <log level="custom">
        <property name="Extracted_Message_ID" expression="$ctx:msgId"/>
        <property name="Extracted_Priority" expression="$ctx:msgPriority"/>
        <property name="Extracted_Queue_Name" expression="$ctx:queueName"/>
    </log>

    <!-- 此处添加数据库更新逻辑 -->
</sequence>

三、更新数据库消息状态

使用DB Mediator执行更新SQL,将提取的元数据作为参数传入。首先在deployment.toml中配置数据源(以MySQL为例):

[database.datasources]
[database.datasources.WSO2_MESSAGE_DB]
url = "jdbc:mysql://localhost:3306/message_db?useSSL=false"
username = "db_user"
password = "db_pass"
driver_class_name = "com.mysql.cj.jdbc.Driver"

接着在序列中添加DB Mediator执行更新:

<!-- 接上述序列代码 -->
<db datasource="WSO2_MESSAGE_DB">
    <statement>
        <sql>UPDATE message_status SET status = 'CONSUMED' WHERE message_id = ? AND queue_name = ?</sql>
        <parameter expression="$ctx:msgId" type="VARCHAR"/>
        <parameter expression="$ctx:queueName" type="VARCHAR"/>
        <!-- 若需用优先级作为条件,可添加对应参数 -->
        <!-- <parameter expression="$ctx:msgPriority" type="INTEGER"/> -->
    </statement>
</db>

<!-- 确认处理完成,若开启手动ACK需添加ACK逻辑 -->
<property name="SET_ROLLBACK_ONLY" value="false" scope="axis2"/>

注意事项

  • 依赖检查:WSO2 MI 4.2默认包含RabbitMQ客户端,若遇兼容问题,可将对应amqp-client.jar放入MI_HOME/lib目录。
  • 消息优先级:RabbitMQ队列需声明时启用优先级(配置x-max-priority参数),否则rabbitmq.message.priority会返回null。
  • 手动ACK:若需手动控制消息确认,可在入站端点添加<parameter name="rabbitmq.auto.ack">false</parameter>,处理完成后添加<property name="rabbitmq.acknowledge.mode" value="ACK" scope="axis2"/>。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 07:20:02