如何为UDP客户端设置多服务器接收优先级?
UDP主备服务器接收优先级处理方案
核心思路
通过跟踪主备服务器的心跳时间戳,结合预设的超时规则,实现数据接收的优先级过滤和主备切换,同时避免频繁跳转。完全匹配需求规则:
- 主服务器活跃时,仅处理主服务器的数据
- 主服务器超时后,切换处理备服务器数据
- 主服务器恢复后,需等备服务器超时才切回(避免频繁跳转)
- 两者均超时触发告警
关键实现步骤
1. 定义状态跟踪变量
在客户端类中维护核心状态:
- 主/备服务器的预设IP地址
- 主/备服务器最后一次收到数据的时间戳
- 超时阈值(可根据实际广播频率调整,示例设为5秒)
- 当前活跃的目标服务器标识(
primary/secondary/none) - 告警触发标记(避免重复告警)
2. 修改接收线程逻辑
在receiver_worker中,收到UDP数据包后先解析来源IP,更新对应服务器的时间戳,再根据当前活跃规则决定是否将数据放入有损队列,直接过滤掉非活跃目标的数据。
3. 定时检查主备状态
启动独立的状态检查线程,定时(示例为每秒一次)校验主备服务器的超时状态,严格按照需求更新当前活跃目标:
- 主服务器未超时:保持/仅在备服务器超时后切回主目标
- 主服务器超时但备服务器活跃:切换到备目标
- 两者均超时:触发告警并标记无活跃目标
完整修改后代码示例
import socket import threading import time from collections import deque # 有损队列实现(按需调整最大长度) class LossyQueue: def __init__(self, max_size=10): self.queue = deque(maxlen=max_size) def put(self, item): self.queue.append(item) def get(self): return self.queue.popleft() if self.queue else None class UDPClient: def __init__(self, listen_port, primary_ip, secondary_ip, timeout=5): self.listen_port = listen_port self.primary_ip = primary_ip self.secondary_ip = secondary_ip self.timeout = timeout # 超时阈值,单位秒 # 状态跟踪变量 self.last_primary_ts = time.time() self.last_secondary_ts = time.time() self.active_target = "primary" # 初始默认主服务器 self.alarm_triggered = False self.lossy_queue = LossyQueue() self.running = True # 启动工作线程 self.receiver_thread = threading.Thread(target=self.receiver_worker, daemon=True) self.check_thread = threading.Thread(target=self.status_check_worker, daemon=True) self.receiver_thread.start() self.check_thread.start() def receiver_worker(self): sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) sock.bind(('', self.listen_port)) sock.settimeout(1) # 设置接收超时,避免线程阻塞 while self.running: try: data, addr = sock.recvfrom(1024) msg = data.decode().strip() sender_ip = addr[0] # 更新对应服务器的心跳时间戳 if sender_ip == self.primary_ip: self.last_primary_ts = time.time() # 仅当前活跃目标为主服务器时,将数据入队 if self.active_target == "primary": self.lossy_queue.put(("primary", msg)) elif sender_ip == self.secondary_ip: self.last_secondary_ts = time.time() # 仅当前活跃目标为备服务器时,将数据入队 if self.active_target == "secondary": self.lossy_queue.put(("secondary", msg)) except socket.timeout: continue except Exception as e: print(f"接收错误: {e}") break sock.close() def status_check_worker(self): while self.running: current_time = time.time() primary_alive = (current_time - self.last_primary_ts) < self.timeout secondary_alive = (current_time - self.last_secondary_ts) < self.timeout # 严格按照需求执行状态切换 if primary_alive: # 主服务器活跃时,仅在备服务器超时的情况下切回 if self.active_target == "secondary" and not secondary_alive: self.active_target = "primary" self.alarm_triggered = False else: if secondary_alive: # 主服务器超时,备服务器活跃则切换到备 if self.active_target != "secondary": self.active_target = "secondary" self.alarm_triggered = False else: # 两者均超时,触发告警 if not self.alarm_triggered: print("告警:主备服务器均未广播数据!") self.alarm_triggered = True self.active_target = "none" time.sleep(1) # 每秒执行一次状态检查 def stop(self): self.running = False self.receiver_thread.join() self.check_thread.join() # 示例用法 if __name__ == "__main__": client = UDPClient( listen_port=5000, primary_ip="192.168.1.100", secondary_ip="192.168.1.101" ) # 模拟前端从队列取数据更新UI try: while True: item = client.lossy_queue.get() if item: print(f"更新UI: 来源[{item[0]}], 状态[{item[1]}]") time.sleep(0.1) except KeyboardInterrupt: client.stop()
代码说明
- LossyQueue:实现简单的有损队列,超过最大长度自动丢弃旧数据,符合状态类数据允许丢失的需求
- receiver_worker:仅处理当前活跃目标服务器的数据,直接过滤非优先级数据,避免前端收到混杂数据
- status_check_worker:独立线程处理状态判断,不影响接收逻辑的性能,同时严格遵循切换规则,避免频繁跳转
- 告警逻辑:通过标记变量避免重复触发告警,减少冗余提示
内容的提问来源于stack exchange,提问作者JRiggles
相关产品推荐
相关产品推荐

