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

Python UDP流传输:休眠时间外的莫名延迟排查与优化

Python UDP流传输延迟问题排查与优化请求

正在开展校园项目,用Python实现UDP流传输:客户端udp_stream.py向服务器udp_stream_server.py发送消息,设定以40条/秒的速率发送800条消息,但实际执行时长超出预期。已完成以下排查操作但问题未解决,恳请分析延迟原因并提供优化指导:

  • 检查代码潜在瓶颈
  • 调整休眠时间以匹配预期消息速率
  • 确认服务器仅执行消息打印与客户端地址输出

客户端代码(udp_stream.py)

from socket import *
import time

def returnTime():
    t = time.localtime()
    current_time = time.strftime("%H:%M:%S", t)
    milliseconds = int((time.time() % 1) * 1000)  # Extract milliseconds
    return f"{current_time}.{milliseconds:03d}"

def log():
    return f"[CLIENT:] {returnTime()} :"

def udp_stream(target_ip, target_port, message_total, message_rate):
    client_socket = socket(AF_INET, SOCK_DGRAM)
    message_number = 10001
    message_size=1470
    sent_messages = 0

    first_sent = log();


    for x in range(message_total):
        sent_messages += 1

        if sent_messages < message_total:
            message = f"{message_number};{'A'*(message_size - len(str(message_number)) -1)}"

            client_socket.sendto(message.encode(), (target_ip, target_port))
            time.sleep(1 / message_rate)
            message_number = message_number + 1



        elif sent_messages == message_total:
            message = f"{message_number};{'A' * (message_size - len(str(message_number)) - 5)}####"

            client_socket.sendto(message.encode(), (target_ip, target_port))

            message_number = message_number + 1
            time.sleep(1 / message_rate)
            last_sent = log()


    print(first_sent, last_sent)


udp_stream('localhost',8080,800,40)

服务器代码(udp_stream_server.py)

import time
def returnTime():
    t = time.localtime()
    current_time = time.strftime("%H:%M:%S", t)
    milliseconds = int((time.time() % 1) * 1000)  # Extract milliseconds
    return f"{current_time}.{milliseconds:03d}"

from socket import *

serverPort = 8080

serverSocket = socket(AF_INET, SOCK_DGRAM)
serverSocket.bind(('', serverPort))

print(f'[SERVER: {returnTime()}]: UDP Server has started listening on port: {serverPort}')

while True:
    # read client's message and remember client's address (IP and port)
    message, clientAddress = serverSocket.recvfrom(1024)
    # Print message and client address
    print (f"[SERVER: {returnTime()}]: Message from client: ",message.decode())
    print (f"[SERVER: {returnTime()}]: Client-IP: ",clientAddress)

延迟原因分析

  1. time.sleep()精度不足:Python的time.sleep()在多数系统下精度仅10-15ms,40条/秒要求每条间隔25ms,加上消息构建、sendto操作的耗时,循环内的延迟会不断累积,最终导致总时长超标。
  2. 服务器接收与打印阻塞:服务器recvfrom(1024)的缓冲区小于消息实际大小(1470字节),会导致消息分段接收,增加处理开销;同时每条消息都触发两次打印IO操作,IO阻塞会拖慢服务器接收速度,间接导致客户端发送缓冲区堆积,影响发送效率。
  3. 客户端串行执行的累积误差:消息构建、发送、休眠完全串行,每一步的微小延迟都会被放大,800条消息的累积误差会非常明显。

优化方案

客户端优化

  1. 改用精确定时策略:放弃依赖time.sleep(),通过时间戳计算目标发送时间,确保每一条消息都在预定时间发送:
def udp_stream(target_ip, target_port, message_total, message_rate):
    client_socket = socket(AF_INET, SOCK_DGRAM)
    message_number = 10001
    message_size = 1470
    interval = 1.0 / message_rate
    start_time = time.time()
    
    first_sent = log()
    for sent_messages in range(1, message_total+1):
        # 计算当前消息的目标发送时间
        target_time = start_time + (sent_messages - 1) * interval
        # 等待到目标时间(用短sleep降低CPU占用)
        while time.time() < target_time:
            time.sleep(0.0001)
        
        # 构建消息
        if sent_messages < message_total:
            message = f"{message_number};{'A'*(message_size - len(str(message_number)) -1)}"
        else:
            message = f"{message_number};{'A' * (message_size - len(str(message_number)) - 5)}####"
        
        client_socket.sendto(message.encode(), (target_ip, target_port))
        message_number += 1
    
    last_sent = log()
    print(first_sent, last_sent)
  1. 预构建所有消息:提前完成所有消息的编码,减少循环内的计算开销:
def udp_stream(target_ip, target_port, message_total, message_rate):
    client_socket = socket(AF_INET, SOCK_DGRAM)
    message_size = 1470
    interval = 1.0 / message_rate
    start_time = time.time()
    
    # 预构建并编码所有消息
    messages = []
    for num in range(10001, 10001 + message_total):
        if num < 10001 + message_total -1:
            msg = f"{num};{'A'*(message_size - len(str(num)) -1)}"
        else:
            msg = f"{num};{'A' * (message_size - len(str(num)) -5)}####"
        messages.append(msg.encode())
    
    first_sent = log()
    for idx, msg in enumerate(messages):
        target_time = start_time + idx * interval
        while time.time() < target_time:
            time.sleep(0.0001)
        client_socket.sendto(msg, (target_ip, target_port))
    
    last_sent = log()
    print(first_sent, last_sent)
  1. 增大发送缓冲区:避免消息在客户端堆积:
client_socket.setsockopt(SOL_SOCKET, SO_SNDBUF, 1024*1024)  # 设置为1MB缓冲区

服务器优化

  1. 调整接收缓冲区大小:确保一次接收完整消息:
message, clientAddress = serverSocket.recvfrom(2048)  # 2048大于消息实际大小1470
  1. 减少打印频率:降低IO阻塞对接收的影响:
count = 0
while True:
    message, clientAddress = serverSocket.recvfrom(2048)
    count +=1
    if count % 10 == 0:
        print(f"[SERVER: {returnTime()}]: Received {count} messages from {clientAddress}")
  1. 分离接收与打印线程:用队列解耦接收和打印操作,避免接收被IO阻塞:
from socket import *
import time
import queue
import threading

def returnTime():
    t = time.localtime()
    current_time = time.strftime("%H:%M:%S", t)
    milliseconds = int((time.time() % 1) * 1000)
    return f"{current_time}.{milliseconds:03d}"

def print_worker(q):
    while True:
        msg, addr = q.get()
        print(f"[SERVER: {returnTime()}]: Message from client: ", msg.decode())
        print(f"[SERVER: {returnTime()}]: Client-IP: ", addr)
        q.task_done()

serverPort = 8080
serverSocket = socket(AF_INET, SOCK_DGRAM)
serverSocket.bind(('', serverPort))
serverSocket.setsockopt(SOL_SOCKET, SO_RCVBUF, 1024*1024)

# 创建队列和后台打印线程
q = queue.Queue()
threading.Thread(target=print_worker, args=(q,), daemon=True).start()

print(f'[SERVER: {returnTime()}]: UDP Server has started listening on port: {serverPort}')

while True:
    message, clientAddress = serverSocket.recvfrom(2048)
    q.put((message, clientAddress))

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 03:22:02