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

如何通过Solace Python API单次调用读取队列中的多条消息?

使用Solace Python API单次读取队列中所有消息的实现方案

核心思路

利用Solace PubSub+ Python API的PersistentMessageReceiver批量接收能力,通过配置批量参数,结合循环接收逻辑,实现单次调用逻辑内读取队列中所有可用消息。

实现步骤与代码示例

1. 初始化连接与持久消息接收器

先创建Solace会话并绑定到已订阅多主题的目标队列:

import time
from solace.messaging.messaging_service import MessagingService
from solace.messaging.resources.queue import Queue
from solace.messaging.receiver.persistent_message_receiver import PersistentMessageReceiver, MessageBatch, AckMode

# Broker连接配置
broker_props = {
    "solace.messaging.transport.host": "tcp://<你的Broker地址>:55555",
    "solace.messaging.service.vpn-name": "<你的VPN名称>",
    "solace.messaging.authentication.scheme.basic.username": "<用户名>",
    "solace.messaging.authentication.scheme.basic.password": "<密码>"
}

# 目标持久化队列(已预先订阅多主题)
target_queue = Queue.durable_exclusive_queue("<队列名称>")

# 初始化并连接消息服务
messaging_service = MessagingService.builder().from_properties(broker_props).build()
messaging_service.connect()

# 创建带批量配置的持久消息接收器
receiver: PersistentMessageReceiver = messaging_service.create_persistent_message_receiver_builder()\
    .with_ack_mode(AckMode.AUTO_ACKNOWLEDGE)  # 自动确认消息,也可选择手动确认
    .with_message_batch_size(max_batch_size=200)  # 单次批量接收的最大消息数
    .with_message_batch_wait_time(max_wait_time_in_millis=300)  # 凑不齐批量时的等待超时时间
    .build(target_queue)

receiver.start()

2. 读取队列中所有可用消息

通过循环调用批量接收方法,直到返回空批次,即可获取当前队列内的全部消息:

all_received_messages = []

try:
    while True:
        # 接收批量消息,超时1秒则返回空批次
        batch: MessageBatch = receiver.receive_batch(timeout_in_millis=1000)
        messages = batch.get_messages()
        if not messages:
            break  # 无更多消息,退出循环
        # 处理并收集消息
        for msg in messages:
            print(f"主题: {msg.get_destination_name()}, 内容: {msg.get_payload_as_string()}")
            all_received_messages.append(msg)
finally:
    # 释放资源
    receiver.stop()
    messaging_service.disconnect()

print(f"共读取到 {len(all_received_messages)} 条消息")

关键配置说明

  • max_batch_size: 单次批量接收的最大消息数,可根据队列消息量调整(比如设为1000)。
  • max_batch_wait_time: 当队列消息数不足批量上限时,等待多久后返回现有消息,避免无限等待。
  • AckMode: 若需严格可靠性,可改用CLIENT_ACKNOWLEDGE,在处理完批次后调用batch.acknowledge()手动确认。

注意事项

  • 确保目标队列是持久化队列,且已正确完成多主题订阅配置。
  • 若队列消息量极大,建议分批次处理,避免内存溢出。
  • 接收器启动后必须调用start(),使用完毕需调用stop()释放连接资源。

内容的提问来源于stack exchange,提问作者Bharat Kumar Kotrike

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 19:42:14