네트웍 프로그램에서 대량의 트래픽이나 클라이언트/서버를 상대로 데이터를 송수신할때 보통 단일쓰레드로 select 함수를 사용해서 처리한다.
리눅스에서는 이것이 보통 합당한 방식인데 왜냐하면 멀티코어 cpu 상에서 cpu 사용률을 보면 각 cpu가 골고루 사용되는것을 볼수 있다. 이것은 이더넷 카드의 인터럽트가 각 cpu 에 골고루 신호를 보내기때문이다.
하지만 윈도우에서 같은 방식으로 싱글쓰레드 select 함수를 사용하면 cpu 1개가 100% 로 치솟고 나머지 cpu 들은 거의 0% 로 놀고있는 현상을 볼수있다. 이것은 이더넷 카드의 인터럽트가 쿼드코어기준으로 4개의 cpu 에 골고루 인터럽트를 보내지않고 1개의 cpu 에게만 집중적으로 보내기 때문이다. 이것을 막기위해 레지스트리 수정등 여러방법을 써봤지만 다 효과가 없었고 오직 IOCP 기술을 사용해서 부하분산을 할 수 있었다.
IO Completion Port 는 일종의 IO 큐로써 등록된 각 소켓 디스크립터들의 이벤트들을 담고있고 초기에 사용자가 만들었던 쓰레드풀에서 각 쓰레드들이 랜덤하게 큐의 이벤트들을 pop 해서 가져가서 처리하는 방식이다.
IO Completion Port 가 일종의 select 함수 역할을 하는 것이고 이벤트에 대한 처리는 각 쓰레드들이 알아서 하는 방식이다. 이렇게 하면 쓰레드들이 골고루 일을 하기때문에 cpu 부하분산을 시킬수 있다.
2011년 10월 18일 화요일
2011년 10월 14일 금요일
IOCP로 UDP 데이터 수신 - how to use IOCP WSARecvFrom
< 변수들 - variables >
// IOCP 컨텍스트 정의
class IOContext {
WSAOVERLAPPED wsaOverlapped;
int socket;
WSABUF wsaBuf;
char buf[10240];
HANDLE hCompletionPort;
HANDLE hCloseEvent;
BOOL bClosed;
IOContext(int sock);
~IOContext();
};
IOContext::IOContext(int sock)
{
socket = sock;
wsaBuf.len = sizeof(buf);
wsaBuf.buf = buf;
SecureZeroMemory((PVOID)&wsaOverlapped, sizeof(WSAOVERLAPPED));
hCloseEvent = CreateEvent(NULL, FALSE, FALSE, NULL);
bClosed = FALSE;
}
IOContext::~IOContext()
{
CloseHandle(hCloseEvent);
}
List<IOContext> contextList; // IOCP 컨텍스트들 링크드 리스트로 관리
HANDLE hCompletionPort;
SYSTEM_INFO systemInfo;
DWORD dwThreadCount;
HANDLE hThreads[16];
< IOCP 초기화 및 쓰레드 생성 - initialize IOCP Handle and create worker threads >
void init()
{
hCompletionPort = CreateIoCompletionPort(INVALID_HANDLE_VALUE, NULL, 0, 0);
GetSystemInfo(&systemInfo);
dwThreadCount = systemInfo.dwNumberOfProcessors*2;
start();
}
void start()
{
for (int i=0; i < dwThreadCount; i++)
{
hThreads[i] = (HANDLE)_beginthreadex(NULL, 0, WorkerThread, this, 0, NULL);
}
}
< IOCP 에 소켓등록 - register socket to IOCP Handle >
void registerSocket(int socketNum)
{
IOContext *context = new IOContext(socketNum);
contextList.insert(context);
context->hCompletionPort = CreateIoCompletionPort((HANDLE)socketNum, hCompletionPort, (ULONG_PTR)context, 0);
int err;
if (!PostQueuedCompletionStatus(hCompletionPort, 0, (ULONG_PTR)context, &context->wsaOverlapped)) {
if ((err = WSAGetLastError()) != WSA_IO_PENDING)
printf("[%s] PostQueuedCompletionStatus error: %d\r\n", err);
}
}
< IOCP 에 소켓등록해제 - unregister socket from IOCP Handle >
void unregisterSocket(int socketNum)
{
IOContext *context = (IOContext *)contextList.search(socketNum);
closesocket(socketNum);
if (WaitForSingleObject(context->hCloseEvent, 1000*5) == WAIT_TIMEOUT) {
printf("CloseEvent Wait Timeout !!!\n");
}
contextList.remove(context);
IOContext *context = (IOContext *)contextList.search(socketNum);
closesocket(socketNum);
if (WaitForSingleObject(context->hCloseEvent, 1000*5) == WAIT_TIMEOUT) {
printf("CloseEvent Wait Timeout !!!\n");
}
contextList.remove(context);
}
< Worker Thread >
unsigned int __stdcall WorkerThread(LPVOID lpParam)
{
BOOL bSuccess = FALSE;
DWORD dwIoSize;
LPOVERLAPPED lpOverlapped = NULL;
IOContext *context;
while (1)
{
bSuccess = GetQueuedCompletionStatus(hCompletionPort, &dwIoSize, (PULONG_PTR)&context, &lpOverlapped, INFINITE);
if (context->bClosed) continue;
if (context->bClosed) continue;
if (!bSuccess) {
err = GetLastError();
if (err == WSA_OPERATION_ABORTED) {
context->bClosed = TRUE;
SetEvent(context->hCloseEvent);
} else if (err == WSAENOTSOCK) {
context->bClosed = TRUE;
SetEvent(context->hCloseEvent);
}
continue;
}
context->bClosed = TRUE;
SetEvent(context->hCloseEvent);
} else if (err == WSAENOTSOCK) {
context->bClosed = TRUE;
SetEvent(context->hCloseEvent);
}
continue;
}
// DO SOMETHING
......
struct sockaddr_in fromAddress;
int len = sizeof(context->buf);
int addressSize = sizeof(fromAddress);
int bytesRead;
DWORD flag = 0;
int err;
context->wsaBuf.len = sizeof(buf);
context->wsaBuf.buf = context->buf;
if (WSARecvFrom(context->socket, &context->wsaBuf, 1, (LPDWORD)&bytesRead, (LPDWORD)&flag, (struct sockaddr*)&fromAddress, (socklen_t *)&addressSize, &context->wsaOverlapped, NULL) == SOCKET_ERROR) {
if ((err=WSAGetLastError()) != WSA_IO_PENDING) {
printf("[%s] WSARecvFrom error:%d, sock:%d, bytesRead:%d\r\n", __FUNCTION__, err, context->socket, bytesRead);
}
if (err == WSAENOTSOCK) { // invalid socket (Socket operation on nonsocket.)
context->bClosed = TRUE;
SetEvent(context->hCloseEvent);
}
if (err == WSAENOTSOCK) { // invalid socket (Socket operation on nonsocket.)
context->bClosed = TRUE;
SetEvent(context->hCloseEvent);
}
}
}
return 0;
return 0;
}
라벨:
socket
2011년 10월 11일 화요일
MySQL 접속권한 주기
mysql> GRANT ALL ON new_aces_server.* TO root@'172.16.1.15' IDENTIFIED BY 'acest';
라벨:
mysql
2011년 9월 30일 금요일
윈도우에서 gettimeofday 함수구현 - windows gettimeofday function
< 상대시간 >
int gettimeofday(struct timeval* tp, int* tz)
{
LARGE_INTEGER tickNow;
static LARGE_INTEGER tickFrequency;
static BOOL tickFrequencySet = FALSE;
if (tickFrequencySet == FALSE) {
QueryPerformanceFrequency(&tickFrequency);
tickFrequencySet = TRUE;
}
QueryPerformanceCounter(&tickNow);
tp->tv_sec = (long) (tickNow.QuadPart / tickFrequency.QuadPart);
tp->tv_usec = (long) (((tickNow.QuadPart % tickFrequency.QuadPart) * 1000000L) / tickFrequency.QuadPart);
return 0;
}
< 절대시간 - unix 타임 >
#include <stdio.h>
#include < time.h >
#include <windows.h>
#if defined(_MSC_VER) || defined(_MSC_EXTENSIONS)
#define DELTA_EPOCH_IN_MICROSECS 11644473600000000Ui64
#else
#define DELTA_EPOCH_IN_MICROSECS 11644473600000000ULL
#endif
struct timezone
{
int tz_minuteswest; /* minutes W of Greenwich */
int tz_dsttime; /* type of dst correction */
};
int GetTimeOfDay(struct timeval *tv, struct timezone *tz)
{
FILETIME ft;
uint64_t tmpres = 0;
static int tzflag;
if (NULL != tv)
{
// system time을 구하기
GetSystemTimeAsFileTime(&ft);
// unsigned 64 bit로 만들기
tmpres |= ft.dwHighDateTime;
tmpres <<= 32;
tmpres |= ft.dwLowDateTime;
// 100nano를 1micro로 변환하기
tmpres /= 10;
// epoch time으로 변환하기
tmpres -= DELTA_EPOCH_IN_MICROSECS;
// sec와 micorsec으로 맞추기
tv->tv_sec = (tmpres / 1000000UL);
tv->tv_usec = (tmpres % 1000000UL);
}
// timezone 처리
if (NULL != tz)
{
if (!tzflag)
{
_tzset();
tzflag++;
}
tz->tz_minuteswest = _timezone / 60;
tz->tz_dsttime = _daylight;
}
return 0;
}
int64_t GetTimeOfDay()
{
struct timeval now;
GetTimeOfDay(&now, NULL);
int64_t millisecond = (int64_t)(now.tv_sec)*1000 + (int64_t)(now.tv_usec)/1000;
return millisecond;
}
int gettimeofday(struct timeval* tp, int* tz)
{
LARGE_INTEGER tickNow;
static LARGE_INTEGER tickFrequency;
static BOOL tickFrequencySet = FALSE;
if (tickFrequencySet == FALSE) {
QueryPerformanceFrequency(&tickFrequency);
tickFrequencySet = TRUE;
}
QueryPerformanceCounter(&tickNow);
tp->tv_sec = (long) (tickNow.QuadPart / tickFrequency.QuadPart);
tp->tv_usec = (long) (((tickNow.QuadPart % tickFrequency.QuadPart) * 1000000L) / tickFrequency.QuadPart);
return 0;
}
< 절대시간 - unix 타임 >
#include <stdio.h>
#include < time.h >
#include <windows.h>
#if defined(_MSC_VER) || defined(_MSC_EXTENSIONS)
#define DELTA_EPOCH_IN_MICROSECS 11644473600000000Ui64
#else
#define DELTA_EPOCH_IN_MICROSECS 11644473600000000ULL
#endif
struct timezone
{
int tz_minuteswest; /* minutes W of Greenwich */
int tz_dsttime; /* type of dst correction */
};
int GetTimeOfDay(struct timeval *tv, struct timezone *tz)
{
FILETIME ft;
uint64_t tmpres = 0;
static int tzflag;
if (NULL != tv)
{
// system time을 구하기
GetSystemTimeAsFileTime(&ft);
// unsigned 64 bit로 만들기
tmpres |= ft.dwHighDateTime;
tmpres <<= 32;
tmpres |= ft.dwLowDateTime;
// 100nano를 1micro로 변환하기
tmpres /= 10;
// epoch time으로 변환하기
tmpres -= DELTA_EPOCH_IN_MICROSECS;
// sec와 micorsec으로 맞추기
tv->tv_sec = (tmpres / 1000000UL);
tv->tv_usec = (tmpres % 1000000UL);
}
// timezone 처리
if (NULL != tz)
{
if (!tzflag)
{
_tzset();
tzflag++;
}
tz->tz_minuteswest = _timezone / 60;
tz->tz_dsttime = _daylight;
}
return 0;
}
int64_t GetTimeOfDay()
{
struct timeval now;
GetTimeOfDay(&now, NULL);
int64_t millisecond = (int64_t)(now.tv_sec)*1000 + (int64_t)(now.tv_usec)/1000;
return millisecond;
}
라벨:
vc++
윈도우에서 writev 함수구현 - windows writev function implementation
#include <winsock2.h>
int writev(int sock, struct iovec *iov, int nvecs)
{
DWORD ret;
if (WSASend(sock, (LPWSABUF)iov, nvecs, &ret, 0, NULL, NULL) == 0) {
return ret;
}
return -1;
}
라벨:
socket
윈도우/리눅스 non-blocking socket 만들기 - making windows/linux non-blocking socket
Boolean makeSocketNonBlocking(int sock) {
#if defined(WIN32) || defined(_WIN32) || defined(IMN_PIM)
unsigned long arg = 1;
return ioctlsocket(sock, FIONBIO, &arg) == 0;
#elif defined(VXWORKS)
int arg = 1;
return ioctl(sock, FIONBIO, (int)&arg) == 0;
#else
int curFlags = fcntl(sock, F_GETFL, 0);
return fcntl(sock, F_SETFL, curFlags|O_NONBLOCK) >= 0;
#endif
}
윈도우는 ioctlsocket, 리눅스는 fcntl 사용
라벨:
socket
2011년 9월 29일 목요일
윈도우 소켓버퍼 사이즈 늘리기 - increasing windows socket buffer size
< 현재 소켓 버퍼 사이즈 읽어오기 >
unsigned getBufferSize(int bufOptName,
int socket) {
unsigned curSize;
SOCKLEN_T sizeSize = sizeof curSize;
if (getsockopt(socket, SOL_SOCKET, bufOptName,
(char*)&curSize, &sizeSize) < 0) {
socketErr("getBufferSize() error: ");
return 0;
}
return curSize;
}
< 소켓 버퍼 사이즈 늘리기 설정 >
#ifndef SOCKLEN_T
#define SOCKLEN_T socklen_t
#endif
unsigned increaseBufferTo( int bufOptName,
int socket, unsigned requestedSize) {
// First, get the current buffer size. If it's already at least
// as big as what we're requesting, do nothing.
unsigned curSize = getBufferSize(bufOptName, socket);
// Next, try to increase the buffer to the requested size,
// or to some smaller size, if that's not possible:
while (requestedSize > curSize) {
SOCKLEN_T sizeSize = sizeof requestedSize;
if (setsockopt(socket, SOL_SOCKET, bufOptName,
(char*)&requestedSize, sizeSize) >= 0) {
// success
return requestedSize;
}
requestedSize = (requestedSize+curSize)/2;
}
return getBufferSize( bufOptName, socket);
}
< 함수 사용 >
// 송신버퍼 사이즈 늘리기
unsigned increaseSendBufferTo(int socket, unsigned requestedSize) {
return increaseBufferTo( SO_SNDBUF, socket, requestedSize);
}
// 수신버퍼 사이즈 늘리기
unsigned increaseReceiveBufferTo( int socket, unsigned requestedSize) {
return increaseBufferTo( SO_RCVBUF, socket, requestedSize);
}
< 함수 사용 >
int ret, size = 1024*1024;
if ((ret=increaseSendBufferTo(fsock, size)) != size)
printf("%s failed to increase send buffer size (size:%d, ret:%d)\r\n", __FUNCTION__, size, ret);
unsigned getBufferSize(int bufOptName,
int socket) {
unsigned curSize;
SOCKLEN_T sizeSize = sizeof curSize;
if (getsockopt(socket, SOL_SOCKET, bufOptName,
(char*)&curSize, &sizeSize) < 0) {
socketErr("getBufferSize() error: ");
return 0;
}
return curSize;
}
< 소켓 버퍼 사이즈 늘리기 설정 >
#ifndef SOCKLEN_T
#define SOCKLEN_T socklen_t
#endif
unsigned increaseBufferTo( int bufOptName,
int socket, unsigned requestedSize) {
// First, get the current buffer size. If it's already at least
// as big as what we're requesting, do nothing.
unsigned curSize = getBufferSize(bufOptName, socket);
// Next, try to increase the buffer to the requested size,
// or to some smaller size, if that's not possible:
while (requestedSize > curSize) {
SOCKLEN_T sizeSize = sizeof requestedSize;
if (setsockopt(socket, SOL_SOCKET, bufOptName,
(char*)&requestedSize, sizeSize) >= 0) {
// success
return requestedSize;
}
requestedSize = (requestedSize+curSize)/2;
}
return getBufferSize( bufOptName, socket);
}
< 함수 사용 >
// 송신버퍼 사이즈 늘리기
unsigned increaseSendBufferTo(int socket, unsigned requestedSize) {
return increaseBufferTo( SO_SNDBUF, socket, requestedSize);
}
// 수신버퍼 사이즈 늘리기
unsigned increaseReceiveBufferTo( int socket, unsigned requestedSize) {
return increaseBufferTo( SO_RCVBUF, socket, requestedSize);
}
< 함수 사용 >
int ret, size = 1024*1024;
if ((ret=increaseSendBufferTo(fsock, size)) != size)
printf("%s failed to increase send buffer size (size:%d, ret:%d)\r\n", __FUNCTION__, size, ret);
*** 버퍼 사이즈를 원하는 만큼 늘릴 수 없으면 윈도우 레지스트리 설정을 추가/수정해야한다.
[HKEY_LOCAL_MACHINE \SYSTEM \CurrentControlSet \Services \Afd \Parameters]
DefaultReceiveWindow = 16384
DefaultSendWindow = 16384
DefaultReceiveWindow = 16384
DefaultSendWindow = 16384
* 리부팅할 필요없고 위 레지스트리 값만 넣어주면된다.
라벨:
socket
피드 구독하기:
글 (Atom)