Многопоточное и асинхронное программирование
https://github.com/jbaldwin/libcoro https://github.com/kelbon/kelcoro
https://github.com/lewissbaker/cppcoro https://github.com/David-Haim/concurrencpp/ https://github.com/facebookexperimental/libunifex
https://isocpp.github.io/CppCoreGuidelines/CppCoreGuidelines#Rcoro-capture
https://colinchcpp.github.io/2023-09-25/11-28-57-239780-coroutine-timeouts-and-cancellation-tokens-in-c/ https://markaicode.com/coroutine-cancellation-patterns/
https://habr.com/ru/articles/326138/
Современные операционные системы и микропроцессоры уже давно поддерживает многозадачность и вместе с тем, каждая из этих задач может выполняться в несколько потоков. Это дает ощутимый прирост производительности вычислений и позволяет лучше масштабировать пользовательские приложения и сервера, но за это приходится платить цену — усложняется разработка программы и ее отладка.
Многопоточное и асинхронное программирование в NewLnag, это способы повышения производительности одного экзмпляра приложения при работе под большой нагрузкой для более полной утилизации вычислительных ресурсов компьютера.
⚠️ НЕ РЕАЛИЗОВАНО: весь описанный здесь раздел многопоточного и асинхронного программирования (сопрограммы/короутины, Awaitable, :Async, :AsyncTask, :AsyncPool, @co_yield/@co_await/@co_return) в текущей версии НЕ реализован. Конструкциям AWAIT/YIELD/WHEN_ALL/WHEN_ANY присвоен Kind=Unimplemented, потоки (:Thread) не реализованы
Так же асинхронное программирование улучшает читабельности кода за счет линейной структуры алгоритма без использования функций обраного вызова или различных вариантов обработчиков событий.
Многопроцессность (multiprocessing), когда каждый процесс (экзмпляр приложения) работает только со своей областью памяти, относится к другому типу взаимодействия нескольких экземпляров приложения и тут не рассматривается.
Выбор архитектры приложения, т.е. многопоточное приложение или написание асинхронного кода зависит от типа вычислительной нагрузки, которая бывает двух видов:
- CPU bound — производительность зависит от нагрузки на процессор, когда он активно работает, например, при выполнении математических расчётов.
- I/O bound — производительность зависит от операций ввода-вывода, когда процессор не работает, а преимуществено находится в ожидании данных, например, при обращении к жесткому диску или ожидании ответа каких-нибудь сетевых API. Другими словами, это нагрузка, которая не зависит от скорости работы процессора.
Многопоточное приложение целесообразно использовать при CPU bound характере вычислительной нагрузки и хорошей масштабируемости алгоритма. Переключение между потоками реализуется средствами ОС и происходит неявно в произвольные моменты времени.
При I/O bound характере нагрузки целесообразно применять концепцию асинхронного программирования, которая реализуется с помощью сопрограмм (короутин). Переключение потока выполнения между сопрограммами (короутинами) происходит не в произвольные моменты времени (как при переключении потоков ОС), а только в моменты ожидания данных у Awaitable объекта (функции).
Awaitable функции это обычные функции, которые возвращают Awaitable объект, и в зависимости от наличия данных, выполнение корутины либо продолжается, если данные были возвращаны, либо выполнение приостанавливается, а управления получает другая корутина.
С точки зрения реализации, Awaitable объект, это спецификация Awaitable шаблона конкретными типом данных, который ожидается при успешном получении данных.
Несколько корутин может выполняться в рамках одного потока операционной системы, что автоматически исключает какие либо гонки между потоками, однако это не исключает запуска асинхронных задач в пуле потоков, что может привести к неконсистентности данных при реализации логики работы приложения.
Поэтому доступ к данным внтури корутин желательно делать с использованием объектов межпотоковой синхронизации.
Для создания потока приложения необходимо расшитить шаблон :Thread или вызвать в отдельно потоке обычныую функцию.
Концепция многопоточного и асинхронного выполнения кода может быть реализована непосредственно или с помощью шаблонных типов данных. Шаблонные потоков и асинхоронных задач применяются как надстройки с типовыми обработчиками ошибок и возможностью кастомизировать тип возвращаемого значения:
- :Thread<@ … @>( … ) применяется для создания тела отдельного потока
- :Task<@ … @>( … ) при создания тела функции для асинхронной задачи
тогда как при непосредствыенной реализации потока или асинхронной корутины обработка ошибок должна быть реализована самостоятельно.
<@ … @>AsyncImpl( … ) := :[ … ]Task<@ … @>( … ){
};
и классы :ThreadPoll() для запуска потоков и :AsyncPool() для запуска асинхронных задач.
:ThreadPoll() принимает awaitable<@ <Any, Error> @> ThreadFunction(...)
:AsyncPool() принимает awaitable<@ <Any, Error> @> AsyncFunction(...)
“[]” YY_TOKEN(AWAIT); “[++]” YY_TOKEN(YIELD); “[++” YY_TOKEN_ONLY(YIELD_BEGIN); “++]” YY_TOKEN_ONLY(YIELD_END); “[]” YY_TOKEN(WHEN_ALL); “[]” YY_TOKEN(WHEN_ANY);
который может выполняться как в рамках своего текущего потока, так и в отдельном потоке или даже пуле потоков.
Тело потока или задачи (функция или короутина) передается в качестве аргумента в конструктор соответствующего класса :Thread() или :Task(), или же базовый класс может быть расширен. Использование базового класса (вместо вызова низкоуровневых функций), позволяет более просто реализовать локальные данных для каждого потока (вместо thread_local переменных) и для уменьшения возможных потенциальных ошибок при обработке исключений в каждом отдельном потоке или корутине, так как не обработанные исключения перехватываются родительским классом, чтобы не произошло краха всего приложения.
в этом же потоке, которая выполняется одновремненно с текущей
Awaitable объект, это спецификация Awaitable шаблона конкретными типом данных, который ожидается при успешном получении данных.
:Query := :Awaitable<@ :String @>();
async_read(stream): Awaitable<@ <:String, :Error> @> := {
# Coroutine as an asynchronous task
@try {
@if(poll(stream)){
# Read and return data
@return :Awaitable(read(stream)); # Return data Awaitable<@ :String @>(read(stream))
} @else {
# No data to read
# Switch to another async task
@return :Awaitable(); # suspend Awaitable<@ :String @>()
}
} @catch (...){
@default @return :Awaitable( @latter ); # return error
}
};
@@ co_await @@ := @@ ** @@; ??????????????????????
async_task(stream) := []() { // Coroutine
str := co_await async_read(stream); # ** async_read(stream);
@print(str);
};
@while( @true ){
:Async( &async_func( Open( Input ) ) );
}
Лямбда функцию или сопрограмму (короутину), можно выполнить как асинхронную задачу, относительно текущего потока выполнения и не зависимо от него.
Но в отличии от потоков операционной системы, переключение между асинхронными задачами происходит явно, методами языка программирования при вызове Awaitable функций, что не требует ресурсов процессора и операционной системы на переключение контекста.
Awaitable функции это обычные функции которые возвращают Awaitable объект, и в зависимости от его значения выполпнение корутины продолжается, если данные были возвращаны, либо выполнение корутины приостанавливается до получения новых данных, а поток управления получает другая корутина в этом же потоке, которая выполняется одновремненно с текущей.
!!!!!!!!!!!!!
Недостатки корутин в C++ https://habr.com/ru/companies/ruvds/articles/755246/
!!!!!!!!!!!!!
Так как асинхронное выполнение корутин происходит в рамках одного потока операционной системы, что автоматически исключает какие либо гонки между потоками, однако не исключает запуска асинхронных задач в пуле потоков, что может привести к неконсистентности данных при реализации логики работы приложения.
Поэтому доступ к данным внтури корутин крайне желательно делать с использованием объектов межпотоковой синхронизации.
Отдельный поток приложения
func() := { slepp(1) };
pure() ::- { slepp(1) };
*thread_func := :Thread( &func );
*thread_pure := :Thread( &pure );
*thread_anon := :Thread( _() := { sleep(1) } );
thread_func.join();
thread_pure.join();
thread_anon.join();
Пул потоков приложения
Создание и удалние потоков операционной системы, это относительно тяжелые (медленные) операции. Чтобы ускорить выполнение большого количества маленьких задач в параллельных потоках приложения можно использовать пул потоков. Данный класс заранее создает заданное количество потоков ОС и динамически назначает один из свободных для выполнения новой задачи по мере их добавления.
@print('Main thread %d', :Thread::get_id());
* pool := :ThreadPool(4); # Maximum 4 threads
* i :Integer = 0;
@while(i < 5) {
# Enqueue tasks for execution
pool.enqueue(
_(task: Integer) := {
@print('Task %d is running on thread %d', $task, :Thread::get_id());
sleep(1);
}
);
i += 1;
};
pool.join();
Main thread 140178994147480
Task 0 is running on thread 140178994148928
Task 1 is running on thread 140178985756224
Task 2 is running on thread 140179010934336
Task 3 is running on thread 140179002541632
Task 4 is running on thread 140178994148928
<@ T @> TypeFunc(arg) := { ... };
:<@ T @> TypeClass(arg: T) := Class(){
};
Senders/Receivers в C++26: от теории к практике
https://habr.com/ru/articles/904134/
https://habr.com/ru/companies/sberbank/articles/829098/
Выбор между процессами, потоками или применением асинхронного подхода зависит в первую очередь от нагрузки. Она бывает двух видов:
- CPU bound — нагрузка на процессор, когда он активно работает, например, при выполнении математических расчётов или вычислениях в «тяжёлых» компьютерных играх.
- I/O bound — процессор ожидает операции ввода-вывода, не слишком интенсивно работая. Из примеров можно привести запросы к базам данных или API каких-нибудь сервисов, то есть к внешним ресурсам. Другими словами это нагрузка, длительность обработки которой не зависит от скорости работы процессора.
- Описание задач декларативно
- Остановка или отмена уже запущеной задачи
- Интеграция задач с планировщиками и I/O (указание, как и где выполнять задачу или привязывать её к событиям ввода-вывода)
- Обработка состояний задачи (ошибки, исключения, отмена, сбои)
- Передача результата выполнения одной задачи в другую без вложенных вызовов и бойлерплейта
- Scheduler описывает контекст исполнения (thread-pool, inline, I/O-реактор) без жёсткой привязки к конкретной реализации. Наглядные реализации:
- inline_scheduler — выполняет задачи синхронно сразу,
- static_thread_pool::scheduler из NVIDIA stdexec,
- reactor_scheduler (I/O на io_uring),
- по хабам разбросаны и другие пользовательские GPU/GUI schedulers.
- Задача - Для многопоточноо и асинхронного выполнения кода используются отдельные типы данных, которые принимают одним из аргументов обычно анонимную функцию или короутину и выполняют её параллельно в отдельном потоке или асинхронно.
var := 1;
var: Any := 1;
var:<:Int, :Double> := 1;
var:<:Int + :Double> := 1;
var:<:Ariphmetic - :Complex> := 1;
[[noreturn]] void rethrow_exception( std::exception_ptr p );
Backpressure - можно вызывать set_stopped при переполнении буфера, downstream возобновляет работу через set_value по событию освобождения. Fan-out / Fan-in - динамическое порождение senders и агрегирование через when_all_range для адаптивных задач. Split / Multicast - один источник данных — и несколько потребителей должны получить одинаковый результат, без повторного вычисления.
Retry с экспоненциальным backoff - если задачу бросила ошибка (например, сеть упала), хотим повторить её несколько раз с увеличивающейся задержкой
Batching, Circuit-breaker, throttling, debounce и т.д. и т.д. и т.д.
- Sender — ленивое описание асинхронной операции. Он не начинает работу, пока не будет подписан или не передан адаптеру: Single-value senders: sync_wait, just, then, co_await sender. Multi-value senders: when_all, when_any, when_all_range.
Как это работает? За кулисами Senders/Receivers определены два низкоуровневых CPO (равно как и в coroutine interop (“P3109R0”)): std::execution::connect(Sender, Receiver) — связывает sender и receiver, создавая «operation state» (состояние операции) без запуска задачи std::execution::start(OperationState&) — запускает уже сконструированную операцию, переводя её в активный режим исполнения
Для многопоточноо и асинхронного выполнения кода используются отдельные типы данных, которые принимают одним из аргументов обычно анонимную функцию или короутину и выполняют её параллельно в отдельном потоке или асинхронно.
Класс :Thread() применяется для создания отдельного потока и класс :Task() для создания короутины для асинхронного выполнения кода, который может выполняться как в рамках своего текущего потока, так и в отдельном потоке или даже пуле потоков.
Тело потока или задачи (функция или короутина) передается в качестве аргумента в конструктор соответствующего класса :Thread() или :Task(), или же базовый класс может быть расширен. Использование базового класса (вместо вызова низкоуровневых функций), позволяет более просто реализовать локальные данных для каждого потока (вместо thread_local перемнных) и для уменьшения возможных потенциальных ошибок при обработке исключений в каждом отдельном потоке или корутине, так как не обработанные исключения перехватываются родительским классом, чтобы не произошло краха всего приложения.
Отдельный поток приложения
Любую функцию можно выполнить как отдельный поток приложения, выполнение которого будет происходить не зависимо от основонго потока. Переключения потоков реализуется средствами ОС и происходит неявно в произвольные моменты времени.
func() := { slepp(1) };
pure() ::- { slepp(1) };
*thread_func := :Thread( &func );
*thread_pure := :Thread( &pure );
*thread_anon := :Thread( _() := { sleep(1) } );
thread_func.join();
thread_pure.join();
thread_anon.join();
Пул потоков приложения
Создание и удалние потоков операционной системы, это относительно тяжелые (медленные) операции. Чтобы ускорить выполнение большого количества маленьких задач в параллельных потоках приложения можно использовать пул потоков. Данный класс заранее создает заданное количество потоков ОС и динамически назначает один из свободных для выполнения новой задачи по мере их добавления.
@print('Main thread %d', :Thread::get_id());
* pool := :ThreadPool(4); # Maximum 4 threads
* i :Integer = 0;
@while(i < 5) {
# Enqueue tasks for execution
pool.enqueue(
_(task: Integer) := {
@print('Task %d is running on thread %d', $task, :Thread::get_id());
sleep(1);
}
);
i += 1;
};
pool.join();
Main thread 140178994147480
Task 0 is running on thread 140178994148928
Task 1 is running on thread 140178985756224
Task 2 is running on thread 140179010934336
Task 3 is running on thread 140179002541632
Task 4 is running on thread 140178994148928
Асинхронное выполнение
Лямбда функцию или сопрограмму (короутину), можно выполнить как асинхронную задачу, относительно текущего потока выполнения и не зависимо от него.
Но в отличии от потоков операционной системы, переключение между асинхронными задачами происходит явно, методами языка программирования при вызове Awaitable функций, что не требует ресурсов процессора и операционной системы на переключение контекста.
Awaitable функции это обычные функции которые возвращают Awaitable объект, и в зависимости от его значения выполпнение корутины продолжается, если данные были возвращаны, либо выполнение корутины приостанавливается до получения новых данных, а поток управления получает другая корутина в этом же потоке, которая выполняется одновремненно с текущей.
!!!!!!!!!!!!!
Недостатки корутин в C++ https://habr.com/ru/companies/ruvds/articles/755246/
!!!!!!!!!!!!!
Так как асинхронное выполнение корутин происходит в рамках одного потока операционной системы, что автоматически исключает какие либо гонки между потоками, однако не исключает запуска асинхронных задач в пуле потоков, что может привести к неконсистентности данных при реализации логики работы приложения.
Поэтому доступ к данным внтури корутин крайне желательно делать с использованием объектов межпотоковой синхронизации.
async_read(stream): <:String,>:Awaitable := { // <:AwaitState, :String,>: Awaitable
# Coroutine as an asynchronous task
@if(poll(stream)){
# Read and return data
@return <read(stream),>:Awaitable; # Return data
//@return :Awaitable:Awaitable(true, read(stream) );
} @else {
# No data to read
# Switch to another async task
@return <,>:Awaitable(); # suspend
}
};
@@ co_await @@ := @@ ** @@;
async_task(stream) := []() { // Coroutine
str := co_await async_read(stream); # ** async_read(stream);
@print(str);
};
@while( @true ){
:Async( &async_func( Open( Input ) ) );
}
Пул асинхронных задач
Асинхронные задачи выполянются в одном и том же потоке приложения, но при наличии нескольких физических ядер у процессора, асинхронные задачи целесообразно распредять сразу по нескольким потокам, выполняющимся на разных физическия ядрах CPU.
Для этих целей служат два класса :AsyncTask - запуск асинхронных задач в отдельном потоке операционной системы и класс :AsyncPoll - который заранее создает заданное количество потоков ОС и динамически назначает один из них для выполнения новой асинхронной задачи по мере их добавления.
@print('Main thread %d', :Thread::get_id());
* task := :AsyncTask(4); # Maximum 4 async task
* i: Integer = 0;
@while(i < 5) {
# Enqueue async tasks for execution
task.enqueue(
[](task: Integer) {
@print('AsyncTask %d is running on thread %d', $task, :Thread::get_id());
* counter:Int64 := 0;
* start_time:Int64 := time::microseconds();
@while(time::microseconds() - start_time < 1_000_000){
usleep(1);
@co_yield; # -+;
counter += 1;
}
@print('AsyncTask %d in thread %d done and was called %d times!', $task, :Thread::get_id(), $counter);
# @co_return; # ++;
}
);
i += 1;
};
task.join();
Main thread 140178994147480
AsyncTask 0 is running on thread 140178994148924
AsyncTask 1 is running on thread 140178994148924
AsyncTask 2 is running on thread 140178994148924
AsyncTask 3 is running on thread 140178994148924
AsyncTask 4 is running on thread 140178994148924
* pool := :AsyncPool(2, 5); # Maximum 2 threads and 5 task in each
@while(@true) {
pool.enqueue( &async_func( Open( Input ) )
};
pool.join();