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

如何为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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 03:00:59