Common Intermediate Language и системное программирование в Microsoft .NET

Разработка параллельных приложений для ОС Windows

Разбить на страницы
Показывать лекцию целиком

Применение потоков и волокон

При разработке параллельных приложений недостаточно создать несколько параллельных ветвей кода - необходимо обеспечить их согласованное выполнение и своевременный обмен данными. При этом возникает необходимость использования как общих для разных потоков и процессов данных, так и создания собственных локальных данных, недоступных другим параллельным ветвям кода, может потребоваться синхронизация и даже учет особенностей аппаратуры, на которой данный код будет исполняться.

Пулы потоков, порт завершения ввода-вывода

Одна из типовых задач - разработка серверов, обслуживающих асинхронно поступающие запросы. Реализация однопоточного сервера для такой задачи нецелесообразна, во-первых, потому что запросы могут приходить в то время, пока сервер занят выполнением предыдущего, а, во-вторых, потому что такой сервер не сможет эффективно задействовать многопроцессорную систему. Можно, конечно, запускать несколько экземпляров однопоточного сервера - но в этом случае потребуется разработка специального диспетчера, поддерживающего очередь запросов и их распределение по списку доступных экземпляров сервера. Альтернативным решением является разработка многопоточного сервера, создающего по специальному рабочему потоку для обработки каждого запроса. Этот вариант также имеет свои недостатки: создание и удаление потоков требует затрат времени, которые будут иметь место в обработке каждого запроса; сверх того, создание большого числа одновременно выполняющихся потоков приведет к общему снижению производительности (и значительному увеличению времени обработки каждого конкретного запроса).

Эти соображения приводят к решению, получившему название пула потоков (thread pool). Для реализации пула потоков необходимо создание некоторого количества потоков, занятых обслуживанием запросов, и диспетчера с очередью запросов. При наличии необработанных запросов диспетчер находит свободный поток и передает запрос этому потоку; если свободных потоков нет, то диспетчер ожидает освобождения какого-либо из занятых потоков. Такой подход обеспечивает, с одной стороны, малые затраты на управление потоками, с другой - достаточно высокую загрузку процессоров и хорошую масштабируемость приложения.

Обычно имеет смысл ограничивать общее число рабочих потоков либо числом доступных процессоров, либо кратным этому числом. Если потоки занимаются исключительно вычислительной работой, то их число не должно превышать число процессоров. Если потоки, сверх того, проводят некоторое время в состоянии ожидания (например, при выполнении операций ввода-вывода), то число потоков следует увеличить - решение следует принимать, исходя из доли времени, которое поток проводит в состоянии простоя, и из полного времени обработки запроса.

Достаточно типичная рекомендация: ограничивать число потоков удвоенным числом процессоров. В случае вычислительных потоков накладные потери будут достаточно малы; в случае потоков, занятых вводом-выводом, утилизация процессоров будет близка к полной. Предполагается, что потоки, не занятые вводом-выводом и при этом проводящие много времени в состоянии ожидания, встречаются весьма редко.

Реализация пула потоков является, на самом деле, нетривиальной задачей - необходимо поддерживать очередь запросов и учитывать состояния потоков из пула (поток простаивает; поток выполняется; поток выполняется, но находится в состоянии ожидания). Также следует учитывать возможность выгрузки части данных в файл подкачки (например, выгрузка стека давно не используемого потока) - иногда быстрее подождать завершения работающего потока, чем активировать простаивающий. Для учета всех этих соображений необходимо реализовать поддержку пулов потоков ядром операционной системы, так как на уровне приложения некоторые нужные сведения просто недоступны.

В Windows такая поддержка реализована в виде порта завершения ввода-вывода. Этот объект ядра берет на себя функциональность, необходимую для организации очереди запросов (используя для этого очередь APC) и списков рабочих потоков, обеспечивая оптимальное управление пулом.

С точки зрения разработчика приложения необходимо:

  • создать порт завершения ввода-вывода;
  • создать пул потоков, ожидающий поступления запросов от этого порта;
  • обеспечить передачу запросов порту.
  • Порт завершения создается с помощью функции

    HANDLE CreateIoCompletionPort(
      HANDLE FileHandle, HANDLE ExistingCompletionPort,
      ULONG_PTR CompletionKey, DWORD NumberOfConcurrentThreads
    );

    Эта функция выполняет две разных операции - во-первых, она создает новый порт завершения, и, во-вторых, она ассоциирует порт с завершением операций ввода-вывода с заданным файлом. Обе эти операции могут быть выполнены одновременно одним вызовом функции, а могут быть исполнены раздельно. Более того, вторая операция - ассоциирование порта завершения ввода-вывода с реальным файлом - может вообще не выполняться. Две типичных формы применения функции CreateIoCompletionPort:

  • Создание нового порта завершения ввода-вывода:
    #define CONCURRENTS  4
    
    HANDLE   hCP;
    hCP = CreateIoCompletionPort(
      INVALID_HANDLE_VALUE, NULL, NULL, CONCURRENTS
    );

    При простом создании порта завершения ввода-вывода достаточно указать только максимальное число одновременно работающих потоков (здесь CONCURRENTS, целесообразно ограничивать числом доступных процессоров). Далее, когда будет создаваться пул потоков, в нем можно будет создать и большее число потоков, чем указано при создании порта - система будет отслеживать, чтобы одновременно исполнялся код не более чем указанного числа потоков. При этом поток, перешедший в состояние ожидания, не считается исполняющимся, так что в случае потоков, проводящих часть времени в режиме ожидания, имеет смысл создавать пул потоков большего размера, чем указано при вызове функции CreateIoCompletionPort.

  • Ассоциирование порта с файлом:
    #define SOME_NUMBER  123
    CreateIoCompletionPort( hFile, hCP, SOME_NUMBER, 0 );

    В этом варианте функция CreateIoCompletionPort не создает нового порта, а возвращает переданный ей описатель уже существующего. Существующий порт завершения ввода-вывода можно связать с несколькими различными файлами одновременно; при этом процедура, обслуживающая завершение ввода-вывода, сможет различать, операция с каким именно файлом поставила в очередь данный запрос, с помощью параметра CompletionKey (здесь SOME_NUMBER ), назначаемого разработчиком. Созданный порт можно не ассоциировать ни с одним файлом - тогда с помощью функции PostQueuedCompletionStatus надо будет помещать в очередь порта запросы, имитирующие завершение ввода-вывода.

  • Рассмотрим небольшой пример (проверка ошибок для упрощения пропущена):

    #include <process.h>
    #define _WIN32_WINNT 0x0500
    #include <windows.h>
    
    #define MAXQUERIES	 15
    #define CONCURENTS	 3
    #define POOLSIZE         5
    
    unsigned _ _stdcall PoolProc( void *arg );
    
    int main( void )
    {
      int     i;
      HANDLE  hcport, hthread[ POOLSIZE ];
      DWORD   temp;
      /* создаем порт завершения ввода-вывода */
      hcport = CreateIoCompletionPort(
        INVALID_HANDLE_VALUE, NULL, NULL, CONCURENTS
      );

    После создания порта надо создать пул потоков. Число потоков в пуле обычно превышает число одновременно работающих потоков, задаваемое при создании порта:

    /* создаем пул потоков */
    for ( i = 0; i < POOLSIZE; i++ ) {
      hthread[i] = (HANDLE)_beginthreadex(
        NULL,0,PoolProc,(void*)hcport,0,(unsigned*)temp
      );
    }

    Если созданный порт ассоциирован с одним или несколькими файлами, то после завершения асинхронных операций ввода-вывода в очереди порта будут размещаться асинхронные запросы, которые система будет направлять для обработки потокам из пула. Однако порт завершения ввода-вывода можно и не связывать с файлами - тогда для размещения запроса можно воспользоваться функцией PostQueuedCompletionStatus, которая размещает запросы в очереди без выполнения реальных операций ввода-вывода.

    /* посылаем несколько запросов в порт */
    for ( i = 0; i < MAXQUERIES; i++ ) {
      PostQueuedCompletionStatus( hcport, 1, i, NULL );
      Sleep( 60 );
    }

    Функция помещает в очередь запросов информацию о "как будто" выполненной операции ввода-вывода, полностью повторяя аргументы - такие как размер переданного блока данных, ключ завершения и указатель на структуру OVERLAPPED, содержащую сведения об операции. Мы можем передавать вместо этих значений произвольные данные. В данном примере, скажем, принято, что значение ключа завершения -1 совместно с длиной переданного блока 0 означает необходимость завершить поток:

    /* для завершения работы посылаем специальные запросы */
    for ( i = 0; i < POOLSIZE; i++ ) {
      PostQueuedCompletionStatus(hcport,0, (ULONG_PTR)-1,NULL);
    }
    /* дожидаемся завершения всех потоков пула 
         и закрываем описатели */
    WaitForMultipleObjects( POOLSIZE, hthread, TRUE, INFINITE );
    for ( i = 0; i < POOLSIZE; i++ ) CloseHandle( hthread[i] );
    CloseHandle( hcport );
    return 0;
    }

    В общем виде поток в пуле реализует цикл с выбором запросов из очереди с помощью функции GetQueuedCompletionStatus. Следует внимательно ознакомиться с описанием этой функции, чтобы грамотно обрабатывать возможные ошибочные ситуации и предусмотреть при необходимости завершение работы потока.

    В данном примере поток в течении 0.3 секунды просто ждет, то есть не исполняется, и порт завершения может передать запросы всем потокам пула, хотя их количество превышает максимальное число одновременно работающих потоков, указанное при создании порта:

    unsigned _ _stdcall PoolProc( void *arg )
    {
      DWORD      	size;
      ULONG_PTR   	key;
      LPOVERLAPPED 	lpov; 
      while (
        GetQueuedCompletionStatus(
          (HANDLE)arg, size, key, lpov, INFINITE
      )) {
        /* проверяем условия завершения цикла */
        if ( !size  key == (ULONG_PTR)-1 ) break;
        Sleep( 300 );
      }
      return 0L;
    }

    Рассмотренный механизм управления пулом потоков весьма эффективен, однако требует некоторого объема ручной работы по созданию порта, по созданию пула потоков и по управлению этими потоками.

    В современных реализациях Windows предусмотрена возможность автоматического создания и управления пулом потоков с помощью функции

    BOOL QueueUserWorkItem(
      LPTHREAD_START_ROUTINE QueryFunction,
      PVOID pContext, ULONG Flags
    );

    Эта функция при необходимости создает пул потоков (число потоков в пуле определяется числом процессоров), создает порт завершения ввода-вывода и размещает в очереди порта запрос. Если нужный порт и пул потоков уже созданы, то она просто размещает новый запрос в очереди порта. При обработке запроса будет вызвана указанная параметром QueryFunction процедура с аргументом pContext:

    DWORD WINAPI QueryFunction( PVOID pContext )
    {
      ...
     return 0L;
    }

    Таким образом, управление пулом потоков сильно упрощается, хотя при этом теряется возможность связывания порта завершения ввода-вывода с конкретными файлами и все запросы должны размещаться в очереди явным вызовом функции QueueUserWorkItem.

    Есть и еще одна особенность у такого способа управления пулом - явного механизма задания числа потоков в пуле не предусмотрено. Однако у разработчика есть возможность управлять этим процессом с помощью последнего параметра функции, содержащего специфичные флаги. Так, с помощью флага WT_EXECUTEDEFAULT запрос будет направлен обычному потоку из пула, флаг WT_EXECUTEINIOTHREAD заставит систему обрабатывать запрос в потоке, который находится в состоянии ожидания оповещения (то есть, надо предусмотреть явные вызовы функции типа SleepEx или WaitForMultipleObjectsEx и т.д.). Флаг WT_EXECUTELONGFUNCTION предназначен для случаев, когда обработка запроса может привести к продолжительному ожиданию - тогда система может увеличить число потоков в пуле:

    #include <process.h>
    #define _WIN32_WINNT 0x0500
    #include <windows.h>
    
    #define MAXQUERIES	15
    #define POOLSIZE 	3
    static LONG    		cnt;
    static HANDLE 		hEvent;
    
    DWORD WINAPI QProc( LPVOID lpData )
    {
      int  r = InterlockedIncrement( cnt );
      Sleep( 300 );
      if ( r >= MAXQUERIES ) SetEvent( hEvent );
      return 0L;
    }
    
    int main( void )
    {
      int  i;
      hEvent = CreateEvent( NULL, TRUE, FALSE, 0 );
      /* в пуле будет не менее POOLSIZE потоков */
      for ( i = 0; i < POOLSIZE; i++ ) {
        QueueUserWorkItem( QProc, NULL, WT_EXECUTELONGFUNCTION );
        Sleep( 60 );
      }
      /* остальные запросы будут распределяться между потоками 
    		    пула,даже если их больше, чем число процессоров */
      for ( ; i < MAXQUERIES; i++ ) {
        QueueUserWorkItem( QProc, NULL, WT_EXECUTEDEFAULT );
        Sleep( 60 );
      }
      /* со временем система может уменьшить число потоков пула */
      /* дожидаемся обработки последнего запроса */
      WaitForSingleObject( hEvent, INFINITE );
      CloseHandle( hEvent );
      return 0;
    }

    Последний пример качественно проще, чем пример с явным созданием порта завершения ввода-вывода, хотя часть возможностей порта завершения при этом не может быть использована.

    Память, локальная для потоков и волокон

    При разработке многопоточных приложений возникает необходимость обеспечивать не только параллельное исполнение кода потоков, но также их взаимодействие - обмен данными, доступ к общим, разделяемым всеми потоками данным и изоляцию некоторых данных одного потока от других.

    Поскольку все потоки разделяют общее адресное пространство процесса, то все они имеют общий и равноправный доступ ко всем данным, хранимым в адресном пространстве. Поэтому для потоков, как правило, не существует проблем с передачей данных друг другу - нужна лишь организация корректного взаимного доступа и изоляция собственных данных от данных других потоков.

    Очень часто для изоляции данных достаточно их размещать в стеке - тогда другие потоки смогут получить к ним доступ либо по явно переданным указателям, либо путем сканирования памяти в поисках стеков других потоков и нужных данных в этих стеках, что относится уже к достаточно трудоемким хакерским технологиям. Однако таким образом трудно организовать постоянное хранение данных, и необходимо постоянно явным образом передавать эти данные (или указатель на них) во все вызываемые процедуры; это достаточно неудобно и не всегда возможно.

    Для решения подобных задач в Windows предусмотрен механизм управления данными, локальными для потока (TLS память, Thread Local Storage). Система предоставляет небольшой специальный блок данных, ассоциированный с каждым потоком. В таком блоке возможно в общем случае хранение произвольных данных, однако, так как размеры этого блока крайне малы, то обычно там размещаются указатели на данные большего объема, выделяемые в приложении для каждого потока; в связи с этим ассоциированную с потоком память можно рассматривать как массив двойных слов или массив указателей.

    ОС Windows предоставляет четыре функции, необходимые для работы с локальной для потока памятью. Функция DWORD TlsAlloc(void) выделяет в ассоциированной с потоком памяти двойное слово, индекс которого возвращается вызвавшей процедуре. Если ассоциированный массив полностью использован, возвращаемое значение будет равно TLS_OUT_OF_INDEXES, что сообщает об ошибке выделения ячейки. Функция TlsFree освобождает выделенную ячейку.

    Если поток выделил некоторую ячейку в ассоциированном массиве, то все потоки данного процесса могут обращаться к ячейке с этим индексом - они получат доступ к ячейкам своих собственных ассоциированных массивов и не смогут узнать или изменить значения, сохраненные в этих ячейках другими потоками. Для доступа к данным зарезервированной ячейки используется функция TlsGetValue, возвращающая значение данной ячейки (в виде указателя, т.к. предполагается, что в ячейках хранятся указатели на некоторые структуры данных) и функция TlsSetValue, изменяющая значение в соответствующей ячейке:

    #include <process.h>
    #include <windows.h>
    
    #define THREADS 18
    static DWORD dwTlsData;
    
    void ProcA( int x )
    {
      int  i;
      int  *iptr = (int*)TlsGetValue( dwTlsData );
      for ( i = 0; i < 100; i++ ) iptr[i] = x;
    }
    
    int ProcB( void )
    { 
      int i, x;
      int *iptr = (int*)TlsGetValue( dwTlsData );
      for ( i = x = 0; i < 100; i++ ) x += iptr[i];
      return x;
    }
    
    unsigned __stdcall ThreadProc( void *param )
    {
      TlsSetValue( dwTlsData, (LPVOID)new int array[100] );
      /* выделенные потоком данные размещены в общей куче,
         используемой всеми потоками, однако указатель на
         эти данные известен только потоку-создателю, так
         как сохраняется в локальной для потока области */
      ProcA( (int)param );
      Sleep( 0 );
      if ( ProcB() != 100*(int)param ) { /* ОШИБКА!!! */ }
      delete[] (int*)TlsGetValue( dwTlsData );
      return 0; 
    }
    int main( void )
    {
      HANDLE    hThread[ THREADS ];
      unsigned  dwThread;
      dwTlsData = TlsAlloc();
      /* создаем новые потоки */
      for ( i = 0; i < THREADS; i++ )
        hThread[i]=(HANDLE)_beginthreadex(
          NULL, 0, ThreadProc, (void*)i, 0, dwThread
        );
      /* дождаться завершения созданных потоков */
      WaitForMultipleObjects( THREADS, hThread, TRUE, INFINITE );
      for ( i = 0; i < THREADS; i++ ) CloseHandle( hThread[i] );
      TlsFree( dwTlsData );
      return 0;
    }

    В приведенном примере в функции main выделяется ячейка в ассоциированном списке, индекс которой сохраняется в глобальной переменной dwTlsData, после чего потоки могут сохранять в этой ячейке свои данные.

    В Visual Studio работа с локальной для потока памятью может быть упрощена при использовании _declspec(thread) при описании переменных. В этом случае компилятор будет размещать эти переменные в специальном сегменте данных ( _TLS ), который будет создаваться библиотекой времени исполнения и ссылки на который будут разрешаться с использованием ассоциированной с потоком памяти. Этот способ во многих случаях предпочтительнее явного управления локальной для потока памятью, так как независимо от числа модулей, использующих такой сегмент, будет задействован только один указатель в ассоциированной памяти (построитель объединит в один большой сегмент все _TLS сегменты модулей).

    #include <process.h>
    #include <windows.h>
    
    #define THREADS 18
    _ _declspec(thread) static int iptr[ 100 ];
    
    void ProcA( int x )
    {
      int  i;
    
      for ( i = 0; i < 100; i++ ) iptr[i] = x;
    }
    int ProcB( void )
    {
      int  i, x;
      for ( i = x = 0; i < 100; i++ ) x += iptr[i];
      return x;
    }
    
    unsigned __stdcall ThreadProc( void *param )
    {
      ProcA( (int)param );
      Sleep( 0 );
      if ( ProcB() != 100*(int)param ) { /* ОШИБКА!!! */ }
      return 0;
    }
    
    int main( void )
    {
      HANDLE    	hThread[THREADS];
      unsigned  	dwThread;
      int     	i;
      /* создаем новые потоки */
      for ( i = 0; i < THREADS; i++ )
        hThread[i] = (HANDLE)_beginthreadex(
          NULL, 0, ThreadProc, (void*)i, 0, dwThread
        );
      /* дождаться завершения созданных потоков */
      WaitForMultipleObjects( THREADS, hThread, TRUE, INFINITE );
      for ( i = 0; i < THREADS; i++ ) CloseHandle( hThread[i] );
      TlsFree( dwTlsData );
      return 0;
    }

    Следует внимательно следить за выделением и освобождением данных, указатели на которые сохраняются в TLS памяти (как в случае явного управления, так и при использовании _ _declspec(thread) ). Могут возникнуть две потенциально ошибочных ситуации:

  • TLS память резервируется в то время, когда уже существуют потоки. Это возможно при явном управлении TLS памятью, и для существующих потоков будут зарезервированы ячейки, но придется предусмотреть специальные меры для их корректной инициализации или для исключения их использования до этого.
  • Все случаи завершения потока. Если TLS память содержит какие-либо указатели, то сама TLS память будет освобождена, а вот те данные, указатели на которые хранились в TLS памяти, - нет. Необходимо специально отслеживать все возможные случаи завершения потоков, включая завершение по ошибке, и принимать меры для освобождения выделенной памяти. При использовании _ _declspec(thread) эта ситуация встречается реже, так как позволяет хранить в _TLS сегментах данные любого фиксированного размера.
  • Следует отметить еще один нюанс, связанный с использованием TLS памяти, волокон и оптимизации. В частных случаях волокна могут исполняться разными потоками - при этом одно и то же волокно должно иметь доступ к TLS памяти именно того потока, в котором оно в данный момент исполняется. А если компилятор генерирует оптимизированный код, то он может разместить указатель на данные TLS памяти в каком-либо регистре или временной переменной, что при переключении волокна на другой поток приведет к ошибке - будет использована TLS память предыдущего потока. Чтобы избежать такой ситуации, компилятору можно указать специальный ключ /GT, отключающий некоторые виды оптимизации при работе с TLS памятью. Это может потребоваться в крайне редких случаях - когда приложение использует несколько волокон, исполняемых в нескольких потоках, и при этом волокна должны использовать TLS память потоков.

    Аналогично TLS памяти, Windows поддерживает память, локальную для волокон, - так называемую FLS память, или Fiber Local Storage. При этом FLS память не зависит от того, какой именно поток выполняет данную нить. Для работы с FLS памятью Windows предоставляет набор функций, аналогичный Tls -функциям, отличие заключается только в функции выделения ячейки FLS памяти:

    DWORD FlsAlloc( PFLS_CALLBACK_FUNCTION lpCallback );
    
    VOID WINAPI FlsCallback( PVOID lpFlsData )
    {
      ...
    }

    Функция отличается от ее аналога TlsAlloc указателем на специальную необязательную процедуру FlsCallback, предоставляемую разработчиком. Эта процедура будет вызвана автоматически при освобождении ячейки FLS памяти (как при завершении волокна, так и при завершении потока или возникновении ошибки), и разработчик может легко предоставить средства для освобождения памяти, указатели на которую были сохранены в ячейках FLS памяти.

    DWORD dwFlsID;
    VOID WINAPI FlsCallback( PVOID lpFlsData )
    {
      /* при завершении волокна или потока память будет освобождена */
      delete[] (int*)lpFlsData;
    }
    void initialize( void )
    {
      dwFlsID = FlsAlloc( FlsCallback );
      ...
    }
    void fiberstart( void )
    {
      FlsSetValue( dwFlsID, new int [ 100 ] );
      /* здесь мы можем не следить за освобождением выделенной памяти */
    }

    Остальные функции для работы с FLS аналогичны Tls-функциям как по описаниям, так и по применению.

    Привязка к процессору и системы с неоднородным доступом к памяти

    ОС Windows предоставляет небольшой набор функций, предназначенных для поддержки систем с неоднородным доступом к памяти (NUMA). К таким функциям относятся средства, обеспечивающие выполнение потоков на конкретных процессорах, и функции, позволяющие получить информацию о структуре NUMA машины. В некоторых случаях привязка потоков к процессорам может преследовать и иные цели, чем поддержка NUMA архитектуры. Так, например, привязка потока к процессору может улучшить использование кэша; на некоторых SMP машинах могут возникать проблемы с использованием таймеров высокого разрешения (опирающихся на счетчики процессоров) и т.д.

    Привязка потоков к процессору задается с помощью специального битового вектора (affinity mask), сохраняемого в целочисленной переменной. Каждый бит этого вектора указывает на возможность исполнения потока на процессоре, номер которого совпадает с номером бита. Таким образом, заданием маски сродства можно ограничить множество процессоров, на которых будет выполняться данный поток. В Windows такие маски назначаются процессу (функции GetProcessAffinityMask и SetProcessAffinityMask ) и потоку (функция SetThreadAffinityMask ). Маска, назначаемая потоку, должна быть подмножеством маски процесса. Помимо ограничения множества процессоров, на которых может исполняться поток, может быть целесообразно назначить потоку самый "удобный" для него процессор (по умолчанию - тот, на котором поток был запущен первый раз). Для этого предназначена функция SetThreadIdealProcessor.

    При использовании NUMA систем следует учитывать, что распределение доступных процессоров по узлам NUMA системы не обязательно последовательное - узлы со смежными номерами могут быть с аппаратной точки зрения весьма удалены друг от друга. Функция GetNumaHighestNodeNumber позволяет определить число NUMA узлов, после чего с помощью обращений к функциям GetNumaProcessorNode, GetNumaNodeProcessorMask и GetNumaAvailableMemoryNode можно определить размещение узлов NUMA системы на процессорах и доступную каждому узлу память.

    Страницы:

    Применение потоков и волокон

    При разработке параллельных приложений недостаточно создать несколько параллельных ветвей кода - необходимо обеспечить их согласованное выполнение и своевременный обмен данными. При этом возникает необходимость использования как общих для разных потоков и процессов данных, так и создания собственных локальных данных, недоступных другим параллельным ветвям кода, может потребоваться синхронизация и даже учет особенностей аппаратуры, на которой данный код будет исполняться.

    Пулы потоков, порт завершения ввода-вывода

    Одна из типовых задач - разработка серверов, обслуживающих асинхронно поступающие запросы. Реализация однопоточного сервера для такой задачи нецелесообразна, во-первых, потому что запросы могут приходить в то время, пока сервер занят выполнением предыдущего, а, во-вторых, потому что такой сервер не сможет эффективно задействовать многопроцессорную систему. Можно, конечно, запускать несколько экземпляров однопоточного сервера - но в этом случае потребуется разработка специального диспетчера, поддерживающего очередь запросов и их распределение по списку доступных экземпляров сервера. Альтернативным решением является разработка многопоточного сервера, создающего по специальному рабочему потоку для обработки каждого запроса. Этот вариант также имеет свои недостатки: создание и удаление потоков требует затрат времени, которые будут иметь место в обработке каждого запроса; сверх того, создание большого числа одновременно выполняющихся потоков приведет к общему снижению производительности (и значительному увеличению времени обработки каждого конкретного запроса).

    Эти соображения приводят к решению, получившему название пула потоков (thread pool). Для реализации пула потоков необходимо создание некоторого количества потоков, занятых обслуживанием запросов, и диспетчера с очередью запросов. При наличии необработанных запросов диспетчер находит свободный поток и передает запрос этому потоку; если свободных потоков нет, то диспетчер ожидает освобождения какого-либо из занятых потоков. Такой подход обеспечивает, с одной стороны, малые затраты на управление потоками, с другой - достаточно высокую загрузку процессоров и хорошую масштабируемость приложения.

    Обычно имеет смысл ограничивать общее число рабочих потоков либо числом доступных процессоров, либо кратным этому числом. Если потоки занимаются исключительно вычислительной работой, то их число не должно превышать число процессоров. Если потоки, сверх того, проводят некоторое время в состоянии ожидания (например, при выполнении операций ввода-вывода), то число потоков следует увеличить - решение следует принимать, исходя из доли времени, которое поток проводит в состоянии простоя, и из полного времени обработки запроса.

    Достаточно типичная рекомендация: ограничивать число потоков удвоенным числом процессоров. В случае вычислительных потоков накладные потери будут достаточно малы; в случае потоков, занятых вводом-выводом, утилизация процессоров будет близка к полной. Предполагается, что потоки, не занятые вводом-выводом и при этом проводящие много времени в состоянии ожидания, встречаются весьма редко.

    Реализация пула потоков является, на самом деле, нетривиальной задачей - необходимо поддерживать очередь запросов и учитывать состояния потоков из пула (поток простаивает; поток выполняется; поток выполняется, но находится в состоянии ожидания). Также следует учитывать возможность выгрузки части данных в файл подкачки (например, выгрузка стека давно не используемого потока) - иногда быстрее подождать завершения работающего потока, чем активировать простаивающий. Для учета всех этих соображений необходимо реализовать поддержку пулов потоков ядром операционной системы, так как на уровне приложения некоторые нужные сведения просто недоступны.

    В Windows такая поддержка реализована в виде порта завершения ввода-вывода. Этот объект ядра берет на себя функциональность, необходимую для организации очереди запросов (используя для этого очередь APC) и списков рабочих потоков, обеспечивая оптимальное управление пулом.

    С точки зрения разработчика приложения необходимо:

  • создать порт завершения ввода-вывода;
  • создать пул потоков, ожидающий поступления запросов от этого порта;
  • обеспечить передачу запросов порту.
  • Порт завершения создается с помощью функции

    HANDLE CreateIoCompletionPort(
      HANDLE FileHandle, HANDLE ExistingCompletionPort,
      ULONG_PTR CompletionKey, DWORD NumberOfConcurrentThreads
    );

    Эта функция выполняет две разных операции - во-первых, она создает новый порт завершения, и, во-вторых, она ассоциирует порт с завершением операций ввода-вывода с заданным файлом. Обе эти операции могут быть выполнены одновременно одним вызовом функции, а могут быть исполнены раздельно. Более того, вторая операция - ассоциирование порта завершения ввода-вывода с реальным файлом - может вообще не выполняться. Две типичных формы применения функции CreateIoCompletionPort:

  • Создание нового порта завершения ввода-вывода:
    #define CONCURRENTS  4
    
    HANDLE   hCP;
    hCP = CreateIoCompletionPort(
      INVALID_HANDLE_VALUE, NULL, NULL, CONCURRENTS
    );

    При простом создании порта завершения ввода-вывода достаточно указать только максимальное число одновременно работающих потоков (здесь CONCURRENTS, целесообразно ограничивать числом доступных процессоров). Далее, когда будет создаваться пул потоков, в нем можно будет создать и большее число потоков, чем указано при создании порта - система будет отслеживать, чтобы одновременно исполнялся код не более чем указанного числа потоков. При этом поток, перешедший в состояние ожидания, не считается исполняющимся, так что в случае потоков, проводящих часть времени в режиме ожидания, имеет смысл создавать пул потоков большего размера, чем указано при вызове функции CreateIoCompletionPort.

  • Ассоциирование порта с файлом:
    #define SOME_NUMBER  123
    CreateIoCompletionPort( hFile, hCP, SOME_NUMBER, 0 );

    В этом варианте функция CreateIoCompletionPort не создает нового порта, а возвращает переданный ей описатель уже существующего. Существующий порт завершения ввода-вывода можно связать с несколькими различными файлами одновременно; при этом процедура, обслуживающая завершение ввода-вывода, сможет различать, операция с каким именно файлом поставила в очередь данный запрос, с помощью параметра CompletionKey (здесь SOME_NUMBER ), назначаемого разработчиком. Созданный порт можно не ассоциировать ни с одним файлом - тогда с помощью функции PostQueuedCompletionStatus надо будет помещать в очередь порта запросы, имитирующие завершение ввода-вывода.

  • Рассмотрим небольшой пример (проверка ошибок для упрощения пропущена):

    #include <process.h>
    #define _WIN32_WINNT 0x0500
    #include <windows.h>
    
    #define MAXQUERIES	 15
    #define CONCURENTS	 3
    #define POOLSIZE         5
    
    unsigned _ _stdcall PoolProc( void *arg );
    
    int main( void )
    {
      int     i;
      HANDLE  hcport, hthread[ POOLSIZE ];
      DWORD   temp;
      /* создаем порт завершения ввода-вывода */
      hcport = CreateIoCompletionPort(
        INVALID_HANDLE_VALUE, NULL, NULL, CONCURENTS
      );

    После создания порта надо создать пул потоков. Число потоков в пуле обычно превышает число одновременно работающих потоков, задаваемое при создании порта:

    /* создаем пул потоков */
    for ( i = 0; i < POOLSIZE; i++ ) {
      hthread[i] = (HANDLE)_beginthreadex(
        NULL,0,PoolProc,(void*)hcport,0,(unsigned*)temp
      );
    }

    Если созданный порт ассоциирован с одним или несколькими файлами, то после завершения асинхронных операций ввода-вывода в очереди порта будут размещаться асинхронные запросы, которые система будет направлять для обработки потокам из пула. Однако порт завершения ввода-вывода можно и не связывать с файлами - тогда для размещения запроса можно воспользоваться функцией PostQueuedCompletionStatus, которая размещает запросы в очереди без выполнения реальных операций ввода-вывода.

    /* посылаем несколько запросов в порт */
    for ( i = 0; i < MAXQUERIES; i++ ) {
      PostQueuedCompletionStatus( hcport, 1, i, NULL );
      Sleep( 60 );
    }

    Функция помещает в очередь запросов информацию о "как будто" выполненной операции ввода-вывода, полностью повторяя аргументы - такие как размер переданного блока данных, ключ завершения и указатель на структуру OVERLAPPED, содержащую сведения об операции. Мы можем передавать вместо этих значений произвольные данные. В данном примере, скажем, принято, что значение ключа завершения -1 совместно с длиной переданного блока 0 означает необходимость завершить поток:

    /* для завершения работы посылаем специальные запросы */
    for ( i = 0; i < POOLSIZE; i++ ) {
      PostQueuedCompletionStatus(hcport,0, (ULONG_PTR)-1,NULL);
    }
    /* дожидаемся завершения всех потоков пула 
         и закрываем описатели */
    WaitForMultipleObjects( POOLSIZE, hthread, TRUE, INFINITE );
    for ( i = 0; i < POOLSIZE; i++ ) CloseHandle( hthread[i] );
    CloseHandle( hcport );
    return 0;
    }

    В общем виде поток в пуле реализует цикл с выбором запросов из очереди с помощью функции GetQueuedCompletionStatus. Следует внимательно ознакомиться с описанием этой функции, чтобы грамотно обрабатывать возможные ошибочные ситуации и предусмотреть при необходимости завершение работы потока.

    В данном примере поток в течении 0.3 секунды просто ждет, то есть не исполняется, и порт завершения может передать запросы всем потокам пула, хотя их количество превышает максимальное число одновременно работающих потоков, указанное при создании порта:

    unsigned _ _stdcall PoolProc( void *arg )
    {
      DWORD      	size;
      ULONG_PTR   	key;
      LPOVERLAPPED 	lpov; 
      while (
        GetQueuedCompletionStatus(
          (HANDLE)arg, size, key, lpov, INFINITE
      )) {
        /* проверяем условия завершения цикла */
        if ( !size  key == (ULONG_PTR)-1 ) break;
        Sleep( 300 );
      }
      return 0L;
    }

    Рассмотренный механизм управления пулом потоков весьма эффективен, однако требует некоторого объема ручной работы по созданию порта, по созданию пула потоков и по управлению этими потоками.

    В современных реализациях Windows предусмотрена возможность автоматического создания и управления пулом потоков с помощью функции

    BOOL QueueUserWorkItem(
      LPTHREAD_START_ROUTINE QueryFunction,
      PVOID pContext, ULONG Flags
    );

    Эта функция при необходимости создает пул потоков (число потоков в пуле определяется числом процессоров), создает порт завершения ввода-вывода и размещает в очереди порта запрос. Если нужный порт и пул потоков уже созданы, то она просто размещает новый запрос в очереди порта. При обработке запроса будет вызвана указанная параметром QueryFunction процедура с аргументом pContext:

    DWORD WINAPI QueryFunction( PVOID pContext )
    {
      ...
     return 0L;
    }

    Таким образом, управление пулом потоков сильно упрощается, хотя при этом теряется возможность связывания порта завершения ввода-вывода с конкретными файлами и все запросы должны размещаться в очереди явным вызовом функции QueueUserWorkItem.

    Есть и еще одна особенность у такого способа управления пулом - явного механизма задания числа потоков в пуле не предусмотрено. Однако у разработчика есть возможность управлять этим процессом с помощью последнего параметра функции, содержащего специфичные флаги. Так, с помощью флага WT_EXECUTEDEFAULT запрос будет направлен обычному потоку из пула, флаг WT_EXECUTEINIOTHREAD заставит систему обрабатывать запрос в потоке, который находится в состоянии ожидания оповещения (то есть, надо предусмотреть явные вызовы функции типа SleepEx или WaitForMultipleObjectsEx и т.д.). Флаг WT_EXECUTELONGFUNCTION предназначен для случаев, когда обработка запроса может привести к продолжительному ожиданию - тогда система может увеличить число потоков в пуле:

    #include <process.h>
    #define _WIN32_WINNT 0x0500
    #include <windows.h>
    
    #define MAXQUERIES	15
    #define POOLSIZE 	3
    static LONG    		cnt;
    static HANDLE 		hEvent;
    
    DWORD WINAPI QProc( LPVOID lpData )
    {
      int  r = InterlockedIncrement( cnt );
      Sleep( 300 );
      if ( r >= MAXQUERIES ) SetEvent( hEvent );
      return 0L;
    }
    
    int main( void )
    {
      int  i;
      hEvent = CreateEvent( NULL, TRUE, FALSE, 0 );
      /* в пуле будет не менее POOLSIZE потоков */
      for ( i = 0; i < POOLSIZE; i++ ) {
        QueueUserWorkItem( QProc, NULL, WT_EXECUTELONGFUNCTION );
        Sleep( 60 );
      }
      /* остальные запросы будут распределяться между потоками 
    		    пула,даже если их больше, чем число процессоров */
      for ( ; i < MAXQUERIES; i++ ) {
        QueueUserWorkItem( QProc, NULL, WT_EXECUTEDEFAULT );
        Sleep( 60 );
      }
      /* со временем система может уменьшить число потоков пула */
      /* дожидаемся обработки последнего запроса */
      WaitForSingleObject( hEvent, INFINITE );
      CloseHandle( hEvent );
      return 0;
    }

    Последний пример качественно проще, чем пример с явным созданием порта завершения ввода-вывода, хотя часть возможностей порта завершения при этом не может быть использована.

    Память, локальная для потоков и волокон

    При разработке многопоточных приложений возникает необходимость обеспечивать не только параллельное исполнение кода потоков, но также их взаимодействие - обмен данными, доступ к общим, разделяемым всеми потоками данным и изоляцию некоторых данных одного потока от других.

    Поскольку все потоки разделяют общее адресное пространство процесса, то все они имеют общий и равноправный доступ ко всем данным, хранимым в адресном пространстве. Поэтому для потоков, как правило, не существует проблем с передачей данных друг другу - нужна лишь организация корректного взаимного доступа и изоляция собственных данных от данных других потоков.

    Очень часто для изоляции данных достаточно их размещать в стеке - тогда другие потоки смогут получить к ним доступ либо по явно переданным указателям, либо путем сканирования памяти в поисках стеков других потоков и нужных данных в этих стеках, что относится уже к достаточно трудоемким хакерским технологиям. Однако таким образом трудно организовать постоянное хранение данных, и необходимо постоянно явным образом передавать эти данные (или указатель на них) во все вызываемые процедуры; это достаточно неудобно и не всегда возможно.

    Для решения подобных задач в Windows предусмотрен механизм управления данными, локальными для потока (TLS память, Thread Local Storage). Система предоставляет небольшой специальный блок данных, ассоциированный с каждым потоком. В таком блоке возможно в общем случае хранение произвольных данных, однако, так как размеры этого блока крайне малы, то обычно там размещаются указатели на данные большего объема, выделяемые в приложении для каждого потока; в связи с этим ассоциированную с потоком память можно рассматривать как массив двойных слов или массив указателей.

    ОС Windows предоставляет четыре функции, необходимые для работы с локальной для потока памятью. Функция DWORD TlsAlloc(void) выделяет в ассоциированной с потоком памяти двойное слово, индекс которого возвращается вызвавшей процедуре. Если ассоциированный массив полностью использован, возвращаемое значение будет равно TLS_OUT_OF_INDEXES, что сообщает об ошибке выделения ячейки. Функция TlsFree освобождает выделенную ячейку.

    Если поток выделил некоторую ячейку в ассоциированном массиве, то все потоки данного процесса могут обращаться к ячейке с этим индексом - они получат доступ к ячейкам своих собственных ассоциированных массивов и не смогут узнать или изменить значения, сохраненные в этих ячейках другими потоками. Для доступа к данным зарезервированной ячейки используется функция TlsGetValue, возвращающая значение данной ячейки (в виде указателя, т.к. предполагается, что в ячейках хранятся указатели на некоторые структуры данных) и функция TlsSetValue, изменяющая значение в соответствующей ячейке:

    #include <process.h>
    #include <windows.h>
    
    #define THREADS 18
    static DWORD dwTlsData;
    
    void ProcA( int x )
    {
      int  i;
      int  *iptr = (int*)TlsGetValue( dwTlsData );
      for ( i = 0; i < 100; i++ ) iptr[i] = x;
    }
    
    int ProcB( void )
    { 
      int i, x;
      int *iptr = (int*)TlsGetValue( dwTlsData );
      for ( i = x = 0; i < 100; i++ ) x += iptr[i];
      return x;
    }
    
    unsigned __stdcall ThreadProc( void *param )
    {
      TlsSetValue( dwTlsData, (LPVOID)new int array[100] );
      /* выделенные потоком данные размещены в общей куче,
         используемой всеми потоками, однако указатель на
         эти данные известен только потоку-создателю, так
         как сохраняется в локальной для потока области */
      ProcA( (int)param );
      Sleep( 0 );
      if ( ProcB() != 100*(int)param ) { /* ОШИБКА!!! */ }
      delete[] (int*)TlsGetValue( dwTlsData );
      return 0; 
    }
    int main( void )
    {
      HANDLE    hThread[ THREADS ];
      unsigned  dwThread;
      dwTlsData = TlsAlloc();
      /* создаем новые потоки */
      for ( i = 0; i < THREADS; i++ )
        hThread[i]=(HANDLE)_beginthreadex(
          NULL, 0, ThreadProc, (void*)i, 0, dwThread
        );
      /* дождаться завершения созданных потоков */
      WaitForMultipleObjects( THREADS, hThread, TRUE, INFINITE );
      for ( i = 0; i < THREADS; i++ ) CloseHandle( hThread[i] );
      TlsFree( dwTlsData );
      return 0;
    }

    В приведенном примере в функции main выделяется ячейка в ассоциированном списке, индекс которой сохраняется в глобальной переменной dwTlsData, после чего потоки могут сохранять в этой ячейке свои данные.

    В Visual Studio работа с локальной для потока памятью может быть упрощена при использовании _declspec(thread) при описании переменных. В этом случае компилятор будет размещать эти переменные в специальном сегменте данных ( _TLS ), который будет создаваться библиотекой времени исполнения и ссылки на который будут разрешаться с использованием ассоциированной с потоком памяти. Этот способ во многих случаях предпочтительнее явного управления локальной для потока памятью, так как независимо от числа модулей, использующих такой сегмент, будет задействован только один указатель в ассоциированной памяти (построитель объединит в один большой сегмент все _TLS сегменты модулей).

    #include <process.h>
    #include <windows.h>
    
    #define THREADS 18
    _ _declspec(thread) static int iptr[ 100 ];
    
    void ProcA( int x )
    {
      int  i;
    
      for ( i = 0; i < 100; i++ ) iptr[i] = x;
    }
    int ProcB( void )
    {
      int  i, x;
      for ( i = x = 0; i < 100; i++ ) x += iptr[i];
      return x;
    }
    
    unsigned __stdcall ThreadProc( void *param )
    {
      ProcA( (int)param );
      Sleep( 0 );
      if ( ProcB() != 100*(int)param ) { /* ОШИБКА!!! */ }
      return 0;
    }
    
    int main( void )
    {
      HANDLE    	hThread[THREADS];
      unsigned  	dwThread;
      int     	i;
      /* создаем новые потоки */
      for ( i = 0; i < THREADS; i++ )
        hThread[i] = (HANDLE)_beginthreadex(
          NULL, 0, ThreadProc, (void*)i, 0, dwThread
        );
      /* дождаться завершения созданных потоков */
      WaitForMultipleObjects( THREADS, hThread, TRUE, INFINITE );
      for ( i = 0; i < THREADS; i++ ) CloseHandle( hThread[i] );
      TlsFree( dwTlsData );
      return 0;
    }

    Следует внимательно следить за выделением и освобождением данных, указатели на которые сохраняются в TLS памяти (как в случае явного управления, так и при использовании _ _declspec(thread) ). Могут возникнуть две потенциально ошибочных ситуации:

  • TLS память резервируется в то время, когда уже существуют потоки. Это возможно при явном управлении TLS памятью, и для существующих потоков будут зарезервированы ячейки, но придется предусмотреть специальные меры для их корректной инициализации или для исключения их использования до этого.
  • Все случаи завершения потока. Если TLS память содержит какие-либо указатели, то сама TLS память будет освобождена, а вот те данные, указатели на которые хранились в TLS памяти, - нет. Необходимо специально отслеживать все возможные случаи завершения потоков, включая завершение по ошибке, и принимать меры для освобождения выделенной памяти. При использовании _ _declspec(thread) эта ситуация встречается реже, так как позволяет хранить в _TLS сегментах данные любого фиксированного размера.
  • Следует отметить еще один нюанс, связанный с использованием TLS памяти, волокон и оптимизации. В частных случаях волокна могут исполняться разными потоками - при этом одно и то же волокно должно иметь доступ к TLS памяти именно того потока, в котором оно в данный момент исполняется. А если компилятор генерирует оптимизированный код, то он может разместить указатель на данные TLS памяти в каком-либо регистре или временной переменной, что при переключении волокна на другой поток приведет к ошибке - будет использована TLS память предыдущего потока. Чтобы избежать такой ситуации, компилятору можно указать специальный ключ /GT, отключающий некоторые виды оптимизации при работе с TLS памятью. Это может потребоваться в крайне редких случаях - когда приложение использует несколько волокон, исполняемых в нескольких потоках, и при этом волокна должны использовать TLS память потоков.

    Аналогично TLS памяти, Windows поддерживает память, локальную для волокон, - так называемую FLS память, или Fiber Local Storage. При этом FLS память не зависит от того, какой именно поток выполняет данную нить. Для работы с FLS памятью Windows предоставляет набор функций, аналогичный Tls -функциям, отличие заключается только в функции выделения ячейки FLS памяти:

    DWORD FlsAlloc( PFLS_CALLBACK_FUNCTION lpCallback );
    
    VOID WINAPI FlsCallback( PVOID lpFlsData )
    {
      ...
    }

    Функция отличается от ее аналога TlsAlloc указателем на специальную необязательную процедуру FlsCallback, предоставляемую разработчиком. Эта процедура будет вызвана автоматически при освобождении ячейки FLS памяти (как при завершении волокна, так и при завершении потока или возникновении ошибки), и разработчик может легко предоставить средства для освобождения памяти, указатели на которую были сохранены в ячейках FLS памяти.

    DWORD dwFlsID;
    VOID WINAPI FlsCallback( PVOID lpFlsData )
    {
      /* при завершении волокна или потока память будет освобождена */
      delete[] (int*)lpFlsData;
    }
    void initialize( void )
    {
      dwFlsID = FlsAlloc( FlsCallback );
      ...
    }
    void fiberstart( void )
    {
      FlsSetValue( dwFlsID, new int [ 100 ] );
      /* здесь мы можем не следить за освобождением выделенной памяти */
    }

    Остальные функции для работы с FLS аналогичны Tls-функциям как по описаниям, так и по применению.

    Привязка к процессору и системы с неоднородным доступом к памяти

    ОС Windows предоставляет небольшой набор функций, предназначенных для поддержки систем с неоднородным доступом к памяти (NUMA). К таким функциям относятся средства, обеспечивающие выполнение потоков на конкретных процессорах, и функции, позволяющие получить информацию о структуре NUMA машины. В некоторых случаях привязка потоков к процессорам может преследовать и иные цели, чем поддержка NUMA архитектуры. Так, например, привязка потока к процессору может улучшить использование кэша; на некоторых SMP машинах могут возникать проблемы с использованием таймеров высокого разрешения (опирающихся на счетчики процессоров) и т.д.

    Привязка потоков к процессору задается с помощью специального битового вектора (affinity mask), сохраняемого в целочисленной переменной. Каждый бит этого вектора указывает на возможность исполнения потока на процессоре, номер которого совпадает с номером бита. Таким образом, заданием маски сродства можно ограничить множество процессоров, на которых будет выполняться данный поток. В Windows такие маски назначаются процессу (функции GetProcessAffinityMask и SetProcessAffinityMask ) и потоку (функция SetThreadAffinityMask ). Маска, назначаемая потоку, должна быть подмножеством маски процесса. Помимо ограничения множества процессоров, на которых может исполняться поток, может быть целесообразно назначить потоку самый "удобный" для него процессор (по умолчанию - тот, на котором поток был запущен первый раз). Для этого предназначена функция SetThreadIdealProcessor.

    При использовании NUMA систем следует учитывать, что распределение доступных процессоров по узлам NUMA системы не обязательно последовательное - узлы со смежными номерами могут быть с аппаратной точки зрения весьма удалены друг от друга. Функция GetNumaHighestNodeNumber позволяет определить число NUMA узлов, после чего с помощью обращений к функциям GetNumaProcessorNode, GetNumaNodeProcessorMask и GetNumaAvailableMemoryNode можно определить размещение узлов NUMA системы на процессорах и доступную каждому узлу память.

    Вернуться к учебному плану