如何通过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
相关产品推荐
相关产品推荐

