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

使用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="&gt;">
          <pendingMessageLimitStrategy>
            <constantPendingMessageLimitStrategy limit="1000"/>
          </pendingMessageLimitStrategy>
        </policyEntry>
        <policyEntry gcInactiveDestinations="true" inactiveTimoutBeforeGC="600000" queue="&gt;"/>
      </policyEntries>
    </policyMap>
  </destinationPolicy>
  <plugins>
  </plugins>
  <networkConnectors>
    <networkConnector conduitSubscriptions="false" duplex="true" name="NetworkConnectorTest1" uri="static:(&lt;URI to 2nd broker&gt;)" userName="&lt;user name&gt;"/>
  </networkConnectors>
  <transportConnectors>
    <transportConnector name="openwire" rebalanceClusterClients="true" updateClusterClients="true" updateClusterClientsOnRemove="true"/>
  </transportConnectors>
</broker>

问题根源

  1. 客户端仅配置单个Broker地址,无故障转移备选节点
  2. 未实现连接断开后的自动重连逻辑
  3. 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&amp;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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 12:59:14