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

Python多线程UDP服务器输出不符问题求助

多线程UDP服务器端口显示异常排查

问题描述

编写了基于Python的多线程UDP服务器及客户端代码,运行后输出不符合预期。实际输出中,第二条消息显示来自客户端端口54766,而预期应为端口62564。

实际输出

2024-10-30 10:41:51,003 - INFO - Server listening on 127.0.0.1:65431
Received message: unencrypted from ('127.0.0.1', 62564)
Received message: unencrypted from ('127.0.0.1', 54766)
Message from 127.0.0.1:62564 - hi
**Message from 127.0.0.1:54766 - hello**

预期输出

2024-10-30 10:41:51,003 - INFO - Server listening on 127.0.0.1:65431
Received message: unencrypted from ('127.0.0.1', 62564)
Received message: unencrypted from ('127.0.0.1', 54766)
Message from 127.0.0.1:62564 - hi
**Message from 127.0.0.1:62564 - hello**

服务器代码

import socket
import threading
import logging

# Server settings
HOST = '127.0.0.1'
PORT = 65431

# List to keep track of connected clients
clients = []
clients_lock = threading.Lock()

# Setting up logging
logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')

# Setting up the server
server = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
server.bind((HOST, PORT))
logging.info(f"Server listening on {HOST}:{PORT}")

def broadcast(message, sender_addr):
    with clients_lock:
        for client in clients:
            if client != sender_addr:
                try:
                    server.sendto(message, client)
                except Exception as e:
                    logging.error(f"Error sending message to {client}: {e}")
                    clients.remove(client)

def handle_client(addr):
    while True:
        if len(clients) == 2:
            try:
                message, _ = server.recvfrom(4096)
                if message:
                    decoded_message = message.decode('utf-8', errors='ignore')
                    sender_info = f"{addr[0]}:{addr[1]}"
                    full_message = f"{sender_info}: {decoded_message}"
                    print(f"Message from {sender_info} - {decoded_message}")
                    #broadcast(full_message.encode('utf-8'), addr)
            except UnicodeDecodeError as e:
                logging.error(f"Unicode decode error: {e}")
            except Exception as e:
                logging.error(f"Error handling client message: {e}")
                with clients_lock:
                    if addr in clients:
                        clients.remove(addr)
                break

def main():
    while True:
        with clients_lock:
            if len(clients) < 2:
                # Receive message from client
                message, addr = server.recvfrom(4096)
                decoded_message = message.decode('utf-8', errors='ignore')
                print(f"Received message: {decoded_message} from {addr}")
                
                if addr not in clients:
                    clients.append(addr)
                    threading.Thread(target=handle_client, args=(addr,)).start()

if __name__ == "__main__":
    main()

客户端代码

import socket
import threading

# Server settings
HOST = '127.0.0.1'
PORT = 65431

# Create a UDP socket
client = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)

# Choose encryption method
encryption_method = input("Choose encryption method (unencrypted): ")

# Flag to stop the receiving thread
stop_thread = threading.Event()

def receive_messages(client_socket):
    while not stop_thread.is_set():
        try:
            message, _ = client_socket.recvfrom(4096)
            if encryption_method == 'unencrypted':
                print(f"Server: {message.decode('utf-8', errors='ignore')}")
        except Exception as e:
            if not stop_thread.is_set():
                print(f"Error receiving message: {e}")
            break

# Start a thread to receive messages
thread = threading.Thread(target=receive_messages, args=(client,))
thread.daemon = True
thread.start()

# Send the encryption method choice once
choice_method = encryption_method.encode('utf-8')
client.sendto(choice_method, (HOST, PORT))

try:
    while True:
        sentence = input("")
        if sentence:
            if encryption_method == 'unencrypted':
                message = sentence.encode('utf-8')
                client.sendto(message, (HOST, PORT))
            print(f"You: {sentence}")
except KeyboardInterrupt:
    print("\nConnection closed.")
finally:
    stop_thread.set()
    client.close()

问题原因

核心问题出在服务器的handle_client函数设计上:

  • 多个handle_client线程共享同一个全局UDP套接字server,调用server.recvfrom(4096)时,所有线程都会尝试接收任意客户端的消息,谁先抢到就处理谁的消息。
  • 线程绑定的addr是初始化时的客户端地址,但处理消息时,不管消息实际来自哪个客户端,都会用线程绑定的addr来打印,导致消息归属错误。比如对应端口62564的线程没抢到消息,而对应54766的线程抢到了62564发来的"hello",就会错误显示成来自54766。

修复方案

修改服务器逻辑,让主循环负责接收所有消息,再将消息分发给对应客户端的处理线程,用队列实现消息分发:

import socket
import threading
import logging
from queue import Queue

# Server settings
HOST = '127.0.0.1'
PORT = 65431

# Dictionary to track client queues: addr -> Queue
client_queues = {}
clients_lock = threading.Lock()

# Setting up logging
logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')

# Setting up the server
server = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
server.bind((HOST, PORT))
logging.info(f"Server listening on {HOST}:{PORT}")

def broadcast(message, sender_addr):
    with clients_lock:
        for addr in client_queues:
            if addr != sender_addr:
                try:
                    server.sendto(message, addr)
                except Exception as e:
                    logging.error(f"Error sending message to {addr}: {e}")
                    del client_queues[addr]

def handle_client(addr, queue):
    while True:
        try:
            message = queue.get()
            if message is None:  # Stop signal
                break
            decoded_message = message.decode('utf-8', errors='ignore')
            sender_info = f"{addr[0]}:{addr[1]}"
            print(f"Message from {sender_info} - {decoded_message}")
            #broadcast(f"{sender_info}: {decoded_message}".encode('utf-8'), addr)
        except UnicodeDecodeError as e:
            logging.error(f"Unicode decode error: {e}")
        except Exception as e:
            logging.error(f"Error handling client message: {e}")
            with clients_lock:
                if addr in client_queues:
                    del client_queues[addr]
            break

def main():
    while True:
        # Receive all messages in main loop
        message, addr = server.recvfrom(4096)
        with clients_lock:
            if addr not in client_queues:
                # First message is encryption method
                decoded_message = message.decode('utf-8', errors='ignore')
                print(f"Received message: {decoded_message} from {addr}")
                # Create queue for new client
                q = Queue()
                client_queues[addr] = q
                threading.Thread(target=handle_client, args=(addr, q)).start()
            else:
                # Put message into client's queue
                client_queues[addr].put(message)

if __name__ == "__main__":
    main()

修复逻辑说明

  1. 主循环负责接收所有客户端消息,判断是否为新客户端:
    • 新客户端则创建消息队列,启动对应线程并记录到client_queues字典。
    • 已有客户端则将消息放入对应队列。
  2. 每个handle_client线程只从自己的队列取消息处理,确保消息与客户端地址一一对应,不会出现归属错误。

内容的提问来源于stack exchange,提问作者Audrey_copper

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 18:42:05