如何在Solace中实现消息处理失败时的Unack操作
Solace PubSub+ 1.6.0 Python API 消息Unack与延迟重处理实现
核心方案
在Solace PubSub+ 1.6.0 Python API中,消息处理失败时的Unack及延迟重投递可通过以下方式实现:
- 必须使用持久消费者:临时订阅的消息Unack后会直接丢弃,只有持久订阅的消息才会被重新投递。
- 调用
Message.nack()方法:在处理失败时主动调用该方法,指定重入队列和延迟时间。
代码示例
from solace.messaging.messaging_service import MessagingService from solace.messaging.receiver.message_receiver import MessageHandler from solace.messaging.resources.topic_subscription import TopicSubscription class CustomMessageHandler(MessageHandler): def on_message(self, message): try: # 替换为你的消息处理逻辑 payload = message.get_payload_as_string() print(f"Processing message: {payload}") # 处理成功,确认消息 message.acknowledge() except Exception as e: print(f"处理失败: {str(e)}") # Unack并设置5秒后重新投递 message.nack(requeue=True, delay_seconds=5) # 初始化并连接消息服务 messaging_service = MessagingService.builder() .from_properties({ "solace.messaging.transport.host": "tcp://你的Solace主机:55555", "solace.messaging.service.vpn-name": "你的VPN名称", "solace.messaging.authentication.scheme.basic.username": "用户名", "solace.messaging.authentication.scheme.basic.password": "密码" }) .build() messaging_service.connect() # 创建持久消费者 durable_consumer = messaging_service.create_direct_message_receiver_builder() .with_subscription(TopicSubscription.of("你的订阅主题/*")) .with_durable_name("持久消费者名称") # 必填,标记为持久订阅 .build() # 启动消费者并绑定处理器 durable_consumer.start(CustomMessageHandler()) # 保持程序运行 try: while True: pass except KeyboardInterrupt: durable_consumer.stop() messaging_service.disconnect()
关键说明
- 持久消费者配置:通过
with_durable_name指定持久化名称,确保消费者断开后未确认的消息不会丢失,重连后可继续处理。 nack()参数:requeue=True:将消息重新放回队列,等待后续投递;若设为False,消息会被移至死信队列(需提前配置死信队列)。delay_seconds:设置消息重投递的延迟时间(单位:秒,支持整数或浮点数)。
- 全局延迟配置:如果所有失败消息的延迟时间一致,可在创建消费者时通过
with_retry_delay(延迟秒数)设置全局重投递延迟,无需在每个nack()中单独指定。
额外注意
- 需在Solace控制台配置队列的最大重投递次数,避免消息因处理失败无限循环重投。
- 确保消息处理逻辑的异常捕获覆盖所有可能的失败场景,避免未捕获异常导致消费者崩溃。
内容的提问来源于stack exchange,提问作者user2223296
相关产品推荐
相关产品推荐

