Perl多线程服务器偶发收发消息异常,求正确实现方案
线程化Perl TCP服务器偶发收发失败问题修复
我尝试在Perl线程中实现TCP服务器,服务器有时能正常收发消息,但存在偶发无法收发的情况。以下是我的初始代码、客户端代码以及两次修改版本,需要一个能稳定运行的线程化服务器实现。
初始服务器代码
use strict; use warnings; use utf8; use experimentals; use threads; sub main() { async { use IO::Socket qw(AF_INET SOCK_STREAM SHUT_WR); my $server = IO::Socket->new( Domain => AF_INET, Type => SOCK_STREAM, Proto => 'tcp', LocalHost => '0.0.0.0', LocalPort => 8447, ReusePort => 1, Listen => 10, ) or say "cannot create socket"; say "listening on 8447"; while(1) { my $client = $server->accept(); my $clientAddress = $client->peerhost(); my $clientPort = $client->peerport(); say "Connection from $clientAddress: $clientPort"; my $data = ""; $client->recv($data, 1024); say "recieved data: ", $data, "\n"; $data = "Thread id: " . threads->self->tid(); $client->send($data); $client->close(); } $server->close(); }; } my $monitorThread = async { while(1) { foreach my $thread (threads->list(threads::joinable)) { $thread->detach(); } sleep(1); } }; main(); while(1) { my @threads = threads->list(threads::all); if(scalar @threads == 1) { $monitorThread->detach(); last; } else { sleep(1); } } exit();
客户端代码
use strict; use warnings; use utf8; use experimentals; use IO::Socket qw(AF_INET SOCK_STREAM SHUT_WR); my $client = IO::Socket->new( Domain => AF_INET, Type => SOCK_STREAM, Proto => 'tcp', PeerPort => 8447, PeerHost => '0.0.0.0', ) or say "Cannot create socket: ", $IO::Socket::errstr; my $size = $client->send("Hello World!"); say "Sent data of length: ", $size; my $buffer; $client->recv($buffer, 1024); say "Got back: ", $buffer; $client->close();
修改版本1
use strict; use warnings; use utf8; use experimentals; use threads; $| = 1; sub clientConnection($client) { async { my $clientAddress = $client->peerhost(); my $clientPort = $client->peerport(); say "Connection from $clientAddress: $clientPort"; my $data = ""; $client->recv($data, 1024); say "recieved data: ", $data, "\n"; my $tid = "Thread id: " . threads->self->tid(); $client->send($tid); $client->close(); }; } sub main() { async { use IO::Socket; my $server = IO::Socket->new( Domain => IO::Socket::AF_INET, Type => IO::Socket::SOCK_STREAM, Proto => 'tcp', LocalHost => '0.0.0.0', LocalPort => 8449, ReusePort => 1, Listen => 10, ) or say "cannot create socket"; say "listening on 8449"; while(1) { my $client = $server->accept(); clientConnection($client); } $server->close(); }; } my $signalThread = async { use sigtrap 'handler' => \&signalHandler, qw(INT); }; sub signalHandler($signalName) { say "got an intterupt, print some info"; my @threads = threads->list(threads::all); say "remaining threads: ", scalar @threads; } my $monitorThread = async { while(1) { foreach my $thread (threads->list(threads::joinable)) { $thread->detach(); } sleep(1); } }; main(); while(1) { my @threads = threads->list(threads::all); if(scalar @threads == 1) { $monitorThread->detach(); last; } else { sleep(1); } } exit();
修改版本2(使用Socket模块)
use strict; use warnings; use utf8; use experimentals; use threads; sub clientConnection($client, $newSocket) { async { my ($clientPort, $clientAddr) = unpack_sockaddr_in($client); my $tid = threads->self->tid(); print $newSocket "a message from server $tid"; print "Connection recieved from ", inet_ntoa($clientAddr), ": ", $clientPort , "\n"; close $newSocket; }; } sub main() { async { use Socket; my $port = 8459; my $proto = getprotobyname('tcp'); my $server = "localhost"; socket(SOCKET, AF_INET, SOCK_STREAM, $proto) or say "cannot open socket $!"; setsockopt(SOCKET, SOL_SOCKET, SO_REUSEADDR, 1) or say "cannot set socket option to SO_REUSEADDR $!"; bind(SOCKET, pack_sockaddr_in($port, inet_aton($server))) or say "Cannot bind to port ", $port; listen(SOCKET, 5) or say "listen Error: $!"; say "Server started on port ", $port; my $newSocket; while(my $client = accept($newSocket, SOCKET)) { clientConnection($client, $newSocket); } }; } my $signalThread = async { use sigtrap 'handler' => \&signalHandler, qw(INT); }; sub signalHandler($signalName) { say "got an intterupt, print some info"; my @threads = threads->list(threads::all); say "remaining threads: ", scalar @threads; } my $monitorThread = async { while(1) { foreach my $thread (threads->list(threads::joinable)) { $thread->detach(); } sleep(1); } }; main(); while(1) { my @threads = threads->list(threads::all); if(scalar @threads == 1) { $monitorThread->detach(); last; } else { sleep(1); } } exit();
正确的线程化Perl TCP服务器实现
修复后的服务器代码
use strict; use warnings; use utf8; use threads; use threads::shared; use IO::Socket qw(AF_INET SOCK_STREAM); # 标记服务器是否运行的共享变量 my $running :shared = 1; sub handle_client { my ($client) = @_; my $tid = threads->self->tid(); # 开启自动刷新,避免缓冲导致响应延迟 $client->autoflush(1); my $client_addr = $client->peerhost(); my $client_port = $client->peerport(); print "[$tid] 连接来自 $client_addr:$client_port\n"; # 读取数据并处理错误 my $data = ''; my $bytes_read = $client->recv($data, 1024); if (defined $bytes_read) { chomp $data; print "[$tid] 收到数据: '$data'\n"; # 发送响应 my $response = "来自线程 $tid 的响应: 已收到 '$data'"; $client->send($response); } else { print "[$tid] 读取客户端数据失败: $!\n"; } # 正确关闭连接:先shutdown再close shutdown($client, 2); close($client); print "[$tid] 连接已关闭\n"; } sub server_thread { my $server = IO::Socket->new( Domain => AF_INET, Type => SOCK_STREAM, Proto => 'tcp', LocalHost => '0.0.0.0', LocalPort => 8447, ReuseAddr => 1, # 使用ReuseAddr而非ReusePort,适配单进程线程化场景 Listen => 10, ) or die "无法创建服务器套接字: $!"; print "服务器监听在 0.0.0.0:8447\n"; while ($running) { # 设置服务器套接字为非阻塞,避免退出时卡在accept $server->blocking(0); my $client = $server->accept(); if ($client) { # 创建子线程处理客户端并直接detach,无需额外监控 threads->create(\&handle_client, $client)->detach(); } else { # 无新连接时短暂休眠,降低CPU占用 sleep(0.1); } } close($server); print "服务器已停止\n"; } # 信号处理:优雅退出 $SIG{INT} = sub { print "\n收到中断信号,正在停止服务器...\n"; $running = 0; # 等待所有可join的线程结束 foreach my $thr (threads->list()) { $thr->join() if $thr->is_joinable(); } exit(0); }; # 启动服务器线程并等待其结束 my $server_thr = threads->create(\&server_thread); $server_thr->join();
问题原因与修复说明
- 套接字生命周期管理:原代码中主线程可能在子线程处理前意外释放客户端套接字,修复后让子线程完全接管套接字,确保资源不被提前回收。
- 线程监控逻辑漏洞:原代码的主线程退出条件存在问题,可能导致服务器提前终止,改用共享变量
$running配合信号处理实现优雅启停。 - 错误处理缺失:原代码未处理
recv/send的异常情况,偶发失败可能源于网络波动或客户端异常关闭,添加错误处理后能稳定应对。 - 套接字选项适配:替换
ReusePort为ReuseAddr,更适合单进程线程化服务器的端口复用需求。 - 输出缓冲处理:开启客户端套接字的
autoflush,避免响应因缓冲延迟或丢失。
内容的提问来源于stack exchange,提问作者buchi
相关产品推荐
相关产品推荐

