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

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();

问题原因与修复说明

  1. 套接字生命周期管理:原代码中主线程可能在子线程处理前意外释放客户端套接字,修复后让子线程完全接管套接字,确保资源不被提前回收。
  2. 线程监控逻辑漏洞:原代码的主线程退出条件存在问题,可能导致服务器提前终止,改用共享变量$running配合信号处理实现优雅启停。
  3. 错误处理缺失:原代码未处理recv/send的异常情况,偶发失败可能源于网络波动或客户端异常关闭,添加错误处理后能稳定应对。
  4. 套接字选项适配:替换ReusePort为ReuseAddr,更适合单进程线程化服务器的端口复用需求。
  5. 输出缓冲处理:开启客户端套接字的autoflush,避免响应因缓冲延迟或丢失。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 20:24:55