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

IOCP客户端与服务器数据收发异常问题求助

IOCP客户端实现问题排查

我能实现常规TCP/IP Socket客户端与服务器,但无法支撑上千级并发客户端,因此开始学习IOCP。但IOCP相关资料较少,搜到的大多是服务器代码,难以理解。

用IOCP服务器搭配普通TCP/IP socket的send/recv客户端时,虽能运行但存在数据发送不完整、无法确认发送完成的问题。服务器流量大时,客户端发往服务器的数据会在单次WSARecv中重复,或接收不完整。

现在为客户端实现IOCP后,出现服务器向客户端发送的数据被自身工作线程接收的异常情况。


客户端工作线程代码

DWORD WINAPI CACGameClient::WorkerThread( LPVOID lpParam )
{

int                 nThreadId = (int)lpParam;

    // 这是客户端的工作线程,只有一个线程
    // 如果服务器崩溃或断开连接,就退出线程并启动重连流程(重连逻辑正常)
    // 这是我第一次写IOCP代码,有看不懂的地方可以问我
    
    DbgPrintSuccess( _X( "Thread #%d: Initialized" ), nThreadId );
    
    DWORD               dwBytesTransfered = NULL;
    ULONG_PTR           lpContext = NULL;
    LPOVERLAPPED        lpOverlapped = NULL;
    CClientManager*     pClientManager = NULL;
    PPER_IO_CONTEXT     lpIOContext = NULL;
    
    WSABUF              buffRecv;
    WSABUF              buffSend;
    
    BOOL                nRet = 0;
    
    DWORD               dwSendNumBytes = 0;
    DWORD               dwRecvNumBytes = 0;
    
    while (TRUE)
    {
    
        BOOL bReturn = GetQueuedCompletionStatus( p_client->m_hCompletionPort, &dwBytesTransfered, &lpContext, &lpOverlapped, INFINITE );
        if (!bReturn) {
            DbgPrintError( _X( "Thread #%d: GetQueuedCompletionStatus failed. ( %d )" ), nThreadId, WSAGetLastError() );
            break;
        }
                
        pClientManager = reinterpret_cast<CClientManager*>(lpContext);
        if (!pClientManager) {
            DbgPrintWarning( _X( "Thread #%d: lpContext is NULL." ), nThreadId );
            break;
        }
    
        DWORD dwFlags = 0;
    
        lpIOContext = pClientManager->pIOContext;
        switch (lpIOContext->IOOperation)
        {
        case OP_READ:
        {
            lpIOContext->IOOperation = OP_WRITE;
            lpIOContext->nTotalBytes = dwBytesTransfered;
            lpIOContext->nSentBytes = 0;
            lpIOContext->wsabuf.len = dwBytesTransfered;
    
            nRet = WSASend( p_client->GetSocket(), &lpIOContext->wsabuf, 1,
                &dwSendNumBytes, dwFlags, &(lpIOContext->Overlapped), NULL );
            if (nRet == SOCKET_ERROR && (ERROR_IO_PENDING != WSAGetLastError())) {
                DbgPrintError( "WSASend() failed: %d\n", WSAGetLastError() );
                continue;
            }
                       
                        // 等OP_WRITE完成、WSARecv执行后,会回到这里
                        // 我打算在这里处理数据
            // ProcessData
    
            break;
        }
        case OP_WRITE:
        {
            lpIOContext->IOOperation = OP_READ;
    
            nRet = WSARecv( p_client->GetSocket(), &lpIOContext->wsabuf, 1,
                &dwRecvNumBytes, &dwFlags, &(lpIOContext->Overlapped), NULL );
    
            if (nRet == SOCKET_ERROR && (ERROR_IO_PENDING != WSAGetLastError())) {
                DbgPrintError( "WSARecv() failed: %d\n", WSAGetLastError() );
                continue;
            }
    
                        // 我查了很多代码,发现大家都用操作码
                        // 现在的问题是,我该在哪里获取数据?
                        // WSARecv完成后会触发新的OP_READ,对吧?
                        // 但为什么会回到这个线程?
                
                 }
        break;
        }
    }
    
    std::thread( ReconnectThread, p_client ).detach();
    
    DbgPrintWarning( _X( "Thread #%d: exited" ), nThreadId );
    
    return 0;

}

初始化IO上下文代码

void InitializeContext( IO_OPERATION OpCode )
{
       pIOContext = (PPER_IO_CONTEXT)malloc( sizeof( PER_IO_CONTEXT ) );
   if (pIOContext)
   {
       pIOContext->Overlapped.Internal = 0;
       pIOContext->Overlapped.InternalHigh = 0;
       pIOContext->Overlapped.Offset = 0;
       pIOContext->Overlapped.OffsetHigh = 0;
       pIOContext->Overlapped.hEvent = NULL;
       pIOContext->IOOperation = OpCode;
       pIOContext->pIOContextForward = NULL;
       pIOContext->nTotalBytes = 0;
       pIOContext->nSentBytes = 0;
       pIOContext->wsabuf.buf = pIOContext->Buffer;
       pIOContext->wsabuf.len = sizeof( pIOContext->Buffer );
       ZeroMemory( pIOContext->wsabuf.buf, pIOContext->wsabuf.len );
   }
   else
   {
   DbgPrintError( "HeapAlloc() PER_SOCKET_CONTEXT failed: %d\n", GetLastError() );
   }
}

发送数据到服务器的代码(非完整)

if (u8SendBuffer)
    {
        delete pIOContext;
        pIOContext = NULL;
    
        InitializeContext( OP_WRITE );
    
        pIOContext->wsabuf.buf = (LPSTR)&Info;
        pIOContext->wsabuf.len = (ULONG)u8SendBufCnt;
        pIOContext->nTotalBytes = u8SendBufCnt;
    
        DWORD dwFlag = 0;
        DWORD dwSend = 0;
    
        int nResult = WSASend( GetSocket(), &pIOContext->wsabuf, 1, &dwSend, dwFlag, &(pIOContext->Overlapped), NULL);
        int nLastEr = WSAGetLastError();
    
        if (nResult == NULL)
        {
            if (nLastEr == ERROR_IO_PENDING )
            {
                DbgPrintSuccess( _X( "All buffers has been sent at once. ( %d )" ), dwSend );
            }
            else
            {
                DbgPrintError( _X( "All buffers failed to send at once. ( %d )" ), dwSend );
            }
        }
        else
        {
            DbgPrintWarning( _X( "Sending unsent buffers, dwSent = %d, dwTotalSize = %d, Error = %d " ), dwSend, u8SendBufCnt, nLastEr );
    
            pIOContext->nSentBytes += dwSend;
    
            while (dwSend < u8SendBufCnt)
            {
                int remains = u8SendBufCnt - dwSend;
    
                WSABUF BufRemains;
                BufRemains.buf = u8SendBuffer + dwSend;
                BufRemains.len = remains;
    
                DWORD dwRemains = 0;
                nResult = WSASend( GetSocket(), &BufRemains, 1, &dwRemains, dwFlag, &(pIOContext->Overlapped), NULL );
                nLastEr = WSAGetLastError();
    
                if (nResult == SOCKET_ERROR )
                {
                    if (nLastEr == ERROR_IO_PENDING)
                    {
                        pIOContext->nSentBytes += dwRemains;
    
                        if (dwSend == u8SendBufCnt)
                        {
                            DbgPrintSuccess( _X( "All buffers has been sent. ( %d )" ), dwSend );
                            break;
                        }
                    }
                    else
                    {
                        DbgPrintError( _X( "All buffers failed to send. ( %d )" ), dwSend );
                        break;
                    }
                    
                }
            }
        }
    
        free( u8SendBuffer );
    }

客户端连接服务器的代码(服务器也是IOCP实现)

if (WSAConnect( clientSocket, reinterpret_cast<sockaddr*>(&server_addr), sizeof( server_addr ), NULL, NULL, NULL, NULL ) != SOCKET_ERROR)
{
    SetSocket( clientSocket );

    DbgPrintInfo( _T( "Address: %s-%s" ), server_ip.c_str(), server_port.c_str() );

    if (AssociateCompletionPort( m_hCompletionPort ) )
    {
        bRet = 0;
        CreateThread( NULL, NULL, (LPTHREAD_START_ROUTINE)WorkerThread, NULL, NULL, NULL );
        DWORD dwBytes = 0;
        DWORD dwFlags = 0;

        InitializeContext( OP_WRITE );

        int nBytesRecv = WSASend( GetSocket(), &(pIOContext->wsabuf), 1,
            &dwBytes, dwFlags, &(pIOContext->Overlapped), NULL );

        int lastError = WSAGetLastError();

        if ((SOCKET_ERROR == nBytesRecv) && (WSA_IO_PENDING != lastError))
        {
            DbgPrintError( "Error in Initial Post. (%d)", lastError );
        }
    }
    else
    {
        WSACleanup();
        bRet = 5;
    }
}

之前用非IOCP客户端时,服务器接受连接后执行初始WSARecv偶尔会出现WSA 10054错误。改成IOCP客户端后,我在连接后加了WSASend,想解决这个错误但还没测试,现在核心问题是服务器自收自发数据,同时我对IOCP工作机制理解不足,求帮忙解决。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 07:48:09