使用Python STOMP库对接Amazon MQ ActiveMQ实现故障转移遇阻
Python STOMP连接Amazon MQ ActiveMQ故障转移问题解决
问题场景
在Python应用中使用STOMP库连接Amazon MQ的ActiveMQ Broker,期望实现故障转移:2个Broker(B1、B2)组成集群,生产者P1和消费者C1仅连接B1,当B1故障/重启时,自动切换到B2。但实际B1故障后,P1连接断开,脚本直接终止。
现有代码与配置
生产者代码
# from stomp import Connection import stomp import time host1 = "<host1>" host2 = "<host2>" port1 = 61614 username = "<username>" password = "<password>" addresses = [(host1, port1)] conn = stomp.Connection(addresses) conn.set_ssl( for_hosts=addresses, cert_file="<path to cert file>", key_file="<path to key file>" ) conn.connect(username, password, wait=True) print("Successfully connected.") # conn.send(body="Testing {}".format(int(time.time())), destination="/queue/queue1", headers={"AMQ_SCHEDULED_DELAY": 15000}) import time for i in range(0, 100): conn.send(body="Testing {} {}".format(int(time.time()), i), destination="/queue/queue1", headers={"persistent": "true", "AMQ_SCHEDULED_DELAY": 15000}) time.sleep(3) conn.disconnect()
消费者代码
import stomp import json class MyListener(stomp.ConnectionListener): def on_connecting(self, host_and_port): super().on_connecting(host_and_port) print("connecting to {}".format(host_and_port)) def on_error(self, frame): print("Received an error {}".format(frame.body)) def on_message(self, frame): print("Received a message {}".format(frame.body)) lambda_payload = json.dumps({"body": frame.body}) host1 = "<host 1>" host2 = "<host 2>" port1 = 61614 username = "<username>" password = "<password>" addresses = [(host2, port1)] conn = stomp.Connection(addresses) conn.set_ssl( for_hosts=addresses, cert_file="<path to cert file>", key_file="<path to key file>" ) conn.set_listener("subscriber1", MyListener()) conn.connect(username, password, wait=True) print("Successfully connected") conn.subscribe(destination="/queue/queue1", id=1, ack="auto") while True: pass
B1的ActiveMQ配置
<?xml version="1.0" encoding="UTF-8" standalone="yes"?> <broker schedulePeriodForDestinationPurge="10000" schedulerSupport="true" xmlns="http://activemq.apache.org/schema/core"> <persistenceAdapter> <kahaDB concurrentStoreAndDispatchQueues="false"/> </persistenceAdapter> <destinationPolicy> <policyMap> <policyEntries> <policyEntry gcInactiveDestinations="true" inactiveTimoutBeforeGC="600000" topic=">"> <pendingMessageLimitStrategy> <constantPendingMessageLimitStrategy limit="1000"/> </pendingMessageLimitStrategy> </policyEntry> <policyEntry gcInactiveDestinations="true" inactiveTimoutBeforeGC="600000" queue=">"/> </policyEntries> </policyMap> </destinationPolicy> <plugins> </plugins> <networkConnectors> <networkConnector conduitSubscriptions="false" duplex="true" name="NetworkConnectorTest1" uri="static:(<URI to 2nd broker>)" userName="<user name>"/> </networkConnectors> <transportConnectors> <transportConnector name="openwire" rebalanceClusterClients="true" updateClusterClients="true" updateClusterClientsOnRemove="true"/> </transportConnectors> </broker>
问题根源
- 客户端仅配置单个Broker地址,无故障转移备选节点
- 未实现连接断开后的自动重连逻辑
- ActiveMQ网络集群的
updateClusterClients仅对OpenWire协议生效,STOMP客户端无法接收集群拓扑更新通知
解决措施
1. 配置多Broker地址与自动重连
修改生产者和消费者的Connection初始化,传入所有Broker地址,并启用无限次重连:
addresses = [(host1, port1), (host2, port1)] # reconnect_attempts_max=-1表示无限重试,reconnect_delay为重试间隔(秒) conn = stomp.Connection(host_and_ports=addresses, reconnect_attempts_max=-1, reconnect_delay=5)
2. 实现断开重连的回调逻辑
在ConnectionListener中添加on_disconnected方法,处理重连及重新订阅逻辑:
生产者Listener示例
class ProducerListener(stomp.ConnectionListener): def on_disconnected(self): print("连接断开,尝试重连...") conn.connect(username, password, wait=True) print("重连成功")
消费者Listener修改
def on_disconnected(self): print("与Broker断开连接,启动重连...") conn.connect(username, password, wait=True) print("重连完成,重新订阅队列...") conn.subscribe(destination="/queue/queue1", id=1, ack="auto")
3. 完善ActiveMQ集群配置
- 确保B2的
networkConnectors也配置指向B1的URI,实现双向集群同步 - 确认transportConnector包含STOMP端口配置(Amazon MQ默认STOMP端口为61614):
<transportConnector name="stomp" uri="stomp+ssl://0.0.0.0:61614?maximumConnections=1000&wireFormat.maxFrameSize=104857600"/>
4. 修改后的完整生产者代码
import stomp import time host1 = "<host1>" host2 = "<host2>" port1 = 61614 username = "<username>" password = "<password>" class ProducerListener(stomp.ConnectionListener): def on_disconnected(self): print("连接丢失,尝试重连...") conn.connect(username, password, wait=True) print("重连成功") addresses = [(host1, port1), (host2, port1)] conn = stomp.Connection(host_and_ports=addresses, reconnect_attempts_max=-1, reconnect_delay=5) conn.set_ssl( for_hosts=addresses, cert_file="<path to cert file>", key_file="<path to key file>" ) conn.set_listener("producer", ProducerListener()) conn.connect(username, password, wait=True) print("连接成功") try: for i in range(0, 100): conn.send(body=f"Testing {int(time.time())} {i}", destination="/queue/queue1", headers={"persistent": "true", "AMQ_SCHEDULED_DELAY": 15000}) time.sleep(3) except Exception as e: print(f"发送消息出错: {e}") finally: conn.disconnect()
5. 修改后的完整消费者代码
import stomp import json import time class MyListener(stomp.ConnectionListener): def on_connecting(self, host_and_port): print(f"正在连接 {host_and_port}") def on_error(self, frame): print(f"收到错误: {frame.body}") def on_message(self, frame): print(f"收到消息: {frame.body}") lambda_payload = json.dumps({"body": frame.body}) def on_disconnected(self): print("断开连接,启动重连...") conn.connect(username, password, wait=True) print("重连完成,重新订阅队列...") conn.subscribe(destination="/queue/queue1", id=1, ack="auto") host1 = "<host 1>" host2 = "<host 2>" port1 = 61614 username = "<username>" password = "<password>" addresses = [(host1, port1), (host2, port1)] conn = stomp.Connection(host_and_ports=addresses, reconnect_attempts_max=-1, reconnect_delay=5) conn.set_ssl( for_hosts=addresses, cert_file="<path to cert file>", key_file="<path to key file>" ) conn.set_listener("subscriber1", MyListener()) conn.connect(username, password, wait=True) print("连接成功") conn.subscribe(destination="/queue/queue1", id=1, ack="auto") try: while True: time.sleep(1) except KeyboardInterrupt: print("停止消费者...") conn.disconnect()
内容的提问来源于stack exchange,提问作者ans98
相关产品推荐
相关产品推荐

