Python XMLRPC主从架构下基于Redis实现动态成员协议咨询
基于Redis Pub/Sub实现XMLRPC主从架构节点动态同步方案
核心实现逻辑
不用写死循环轮询get_message(),用redis-py原生的阻塞迭代器+后台守护线程实现事件监听,完全不阻塞主业务逻辑,三类角色(master、业务客户端、后续的存活探测节点)复用同一套逻辑即可:
- 固定Redis Hash键
xmlrpc:worker:registry存全量worker元数据,作为可信的数据源 - 固定Pub/Sub频道
xmlrpc:worker:change广播节点变更事件,事件类型覆盖add(新节点注册)、remove(主动下线)、health_offline(探测宕机下线) - 所有需要感知节点列表的角色启动时先拉取全量列表初始化本地缓存,之后订阅变更频道做增量更新,不用每次发请求前手动拉取列表
- 节点变更操作(注册、下线、踢除宕机节点)执行时,先更新Hash里的可信数据源,再往频道发事件,所有订阅方实时同步
可直接复用的实现代码
先装依赖:pip install redis
公共同步模块代码,所有角色直接实例化调用即可:
import threading import redis import json from typing import Callable, Dict, Optional class DynamicWorkerRegistry: def __init__(self, redis_host: str = "127.0.0.1", redis_port: int = 6379, redis_db: int = 0): # 普通读写用独立连接,和PubSub连接隔离,避免PubSub阻塞导致普通命令失效 self.redis_cli = redis.Redis(host=redis_host, port=redis_port, db=redis_db, decode_responses=True) self.pubsub_cli = self.redis_cli.pubsub() # 本地缓存,业务代码直接读这个变量,永远是最新节点列表 self.worker_list: Dict[str, dict] = {} self.change_cb: Optional[Callable] = None # 固定键和频道名,所有角色统一约定 self._REGISTRY_KEY = "xmlrpc:worker:registry" self._CHANNEL = "xmlrpc:worker:change" self._listen_thread: Optional[threading.Thread] = None self._running = False def _load_full_workers(self): # 启动时一次性拉全量列表初始化 all_data = self.redis_cli.hgetall(self._REGISTRY_KEY) for wid, meta_str in all_data.items(): self.worker_list[wid] = json.loads(meta_str) def _listen_event_loop(self): self.pubsub_cli.subscribe(self._CHANNEL) # 用原生listen()迭代器阻塞等消息,不用自己写while True调get_message for msg in self.pubsub_cli.listen(): if not self._running: break # 过滤掉订阅成功类的系统消息 if msg["type"] != "message": continue try: event = json.loads(msg["data"]) op = event["op"] wid = event["worker_id"] meta = event.get("meta", {}) # 增量更新本地缓存 if op == "add": self.worker_list[wid] = meta elif op in ("remove", "health_offline"): self.worker_list.pop(wid, None) # 触发自定义回调,比如打日志、调整负载均衡权重 if self.change_cb: self.change_cb(op, wid, meta) except Exception as e: print(f"处理节点变更事件异常: {str(e)}") def start(self, on_change: Optional[Callable] = None): """启动同步,初始化时调用一次即可""" if self._running: return self.change_cb = on_change self._load_full_workers() self._running = True # 设为守护线程,主进程退出时自动回收,不用手动处理线程关闭 self._listen_thread = threading.Thread(target=self._listen_event_loop, daemon=True) self._listen_thread.start() def stop(self): self._running = False self.pubsub_cli.unsubscribe(self._CHANNEL) # ========== Master/存活探测节点调用的变更方法 ========== def add_worker(self, worker_id: str, worker_meta: dict): """注册新worker时调用""" self.redis_cli.hset(self._REGISTRY_KEY, worker_id, json.dumps(worker_meta)) event = {"op": "add", "worker_id": worker_id, "meta": worker_meta} self.redis_cli.publish(self._CHANNEL, json.dumps(event)) def del_worker(self, worker_id: str, is_offline_by_healthcheck: bool = False): """删除worker/踢除宕机节点时调用""" self.redis_cli.hdel(self._REGISTRY_KEY, worker_id) op_type = "health_offline" if is_offline_by_healthcheck else "remove" event = {"op": op_type, "worker_id": worker_id} self.redis_cli.publish(self._CHANNEL, json.dumps(event))
不同角色使用示例
Master节点
registry = DynamicWorkerRegistry() # master如果需要感知节点状态也可以调start()启动监听 # 新worker注册 registry.add_worker("worker_01", {"ip": "10.0.0.2", "port": 8000, "max_task": 10}) # 主动移除worker registry.del_worker("worker_01")
业务客户端/存活探测节点
def change_log(op, wid, meta): print(f"节点变更:操作{op},节点{wid},当前可用节点数{len(registry.worker_list)}") registry = DynamicWorkerRegistry() # 启动同步,传入变更回调 registry.start(on_change=change_log) # 业务逻辑直接读registry.worker_list即可,不需要主动拉取 # 比如做XMLRPC调用时直接从本地缓存选节点 import random if registry.worker_list: selected_wid, selected_meta = random.choice(list(registry.worker_list.items())) # 初始化对应XMLRPC客户端发请求即可
注意事项
- 不要把PubSub连接和普通Redis读写连接混用,PubSub会将连接置为阻塞模式,会导致普通Redis命令执行失败,代码里已经做了连接隔离
- 如果需要更高可靠性,可以在监听循环里加连接异常捕获,触发重连即可,不需要改动上层逻辑
- 后续开发存活探测节点时,直接调用
del_worker方法传入is_offline_by_healthcheck=True即可踢除宕机节点,所有客户端会自动同步移除该节点,完全适配动态协议要求 - 整个同步逻辑只有启动时做一次全量拉取,后续全部是事件增量推送,没有轮询开销,比原生Socket组播方案少维护大量连接状态逻辑,稳定性更高
内容的提问来源于stack exchange,提问作者Deividgp
相关产品推荐
相关产品推荐

