Elixir Конкурентность в Elixir: процессы, OTP, супервизия, GenStage/Broadway и распределённость
0%

Конкурентность в Elixir: процессы, OTP, супервизия, GenStage/Broadway и распределённость

Конкурентность в Elixir

Это центральная статья курса. Всё остальное в Elixir — приятно, но именно модель конкурентности отличает его от любого мейнстрим-языка. Здесь вы научитесь мыслить в терминах процессов и деревьев супервизии — навык, который переносится на построение по-настоящему отказоустойчивых систем.

С первых принципов: почему акторы

В большинстве языков конкурентность — это потоки, разделяющие память, и вы защищаете эту память замками (мьютексами). Замки́ трудны: их легко забыть, легко устроить дедлок, легко получить гонку. Отладка таких багов — ад.

BEAM выбирает другую модель — акторную. Единица конкурентности — процесс, который:

  • ничего не разделяет с другими (share-nothing);
  • имеет свой почтовый ящик (mailbox) и обрабатывает сообщения по одному;
  • общается только асинхронными сообщениями.

Раз нет разделяемой памяти — нет и гонок за неё, и мьютексов в прикладном коде становится не нужно. Конкурентность становится безопасной по построению.

Процессы BEAM — не ОС-потоки и не зелёные потоки в привычном смысле. Они запредельно дёшевы (~300 слов на старте), их запускают миллионами, и планировщик BEAM вытесняюще распределяет их по ядрам. Один жадный процесс не заморозит остальные.

Базовые примитивы: spawn, send, receive

Самый низкий уровень (в проде вы будете пользоваться OTP-абстракциями поверх, но понимать основу обязательно):

# spawn создаёт процесс и возвращает его PID (идентификатор процесса)
pid = spawn(fn -> IO.puts("Привет из процесса #{inspect(self())}") end)

# send отправляет сообщение в почтовый ящик (асинхронно, не блокирует)
send(pid, {:hello, "мир"})

# receive забирает сообщение из ящика, сопоставляя по образцу
receive do
  {:hello, msg} -> IO.puts("Получил: #{msg}")
  {:bye, _} -> IO.puts("Пока")
after
  5_000 -> IO.puts("Таймаут — за 5 секунд ничего не пришло")
end

Ключевые свойства:

  • send не блокирует и не гарантирует, что получатель жив — просто кладёт в ящик.
  • receive блокирует процесс, пока не придёт подходящее сообщение (или сработает after).
  • Сообщения обрабатываются по порядку прихода, но receive может выборочно (по паттерну) забирать из ящика.

Классический паттерн — рекурсивный цикл с состоянием (это и есть «сервер» на голых примитивах):

defmodule Counter do
  def start(initial), do: spawn(fn -> loop(initial) end)

  defp loop(state) do
    receive do
      {:inc, n} -> loop(state + n)                # новое состояние = хвостовая рекурсия
      {:get, from} ->
        send(from, {:count, state})
        loop(state)
      :stop -> :ok                                 # выходим из цикла — процесс завершится
    end
  end
end

Именно эту рутину («держи состояние, принимай запросы, отвечай, зациклись») инкапсулирует GenServer. Дальше мы почти всегда используем OTP, а не голый spawn.

Ссылки и мониторы: как узнать о падении

Процессы изолированы, но часто должны знать о судьбе друг друга.

Ссылка (link) — двунаправленная связь: если один процесс падает, сигнал завершения распространяется на связанный, и тот тоже падает (если не «ловит выходы»). Так падения распространяются по группе связанных процессов — это фундамент супервизии.

spawn_link(fn -> raise "бум" end)   # если этот упадёт — уронит и родителя

Монитор (monitor) — однонаправленное наблюдение: наблюдатель получает сообщение о падении наблюдаемого, но сам не падает. Идеально, когда нужно просто узнать о завершении.

{pid, ref} = spawn_monitor(fn -> raise "бум" end)

receive do
  {:DOWN, ^ref, :process, ^pid, reason} ->
    IO.puts("Процесс упал: #{inspect(reason)}")
end

Trap exits — процесс может «ловить» сигналы завершения вместо падения (Process.flag(:trap_exit, true)), получая их как обычные сообщения {:EXIT, pid, reason}. Так устроены супервизоры внутри.

Механизм Направление Реакция на падение
link двунаправленный связанный тоже падает (или ловит :EXIT)
monitor однонаправленный наблюдатель получает :DOWN, живёт дальше

«Let it crash»: философия, а не безрассудство

Раз падения изолированы (свой стек, своя куча) и распространяются контролируемо (link) или наблюдаемо (monitor), появляется мощная стратегия: не программируй защиту от неожиданного — дай процессу упасть и восстанови его в известное хорошее состояние.

Почему это лучше тотального try/catch:

  1. Прикладной код проще — вы пишете happy path, не засоряя его обработкой каждой аномалии.
  2. Известное хорошее состояние после рестарта часто надёжнее, чем попытка «починить» повреждённое состояние на лету.
  3. Баг в одном запросе не роняет систему — падает один процесс (одно соединение, один заказ), остальные работают.

«Let it crash» не означает «не проверять ввод». Ожидаемые ошибки (валидация) обрабатывайте значениями ({:error, _}). Падать нужно на том, что не должно было случиться — на нарушенных инвариантах. Восстановление обеспечивает супервизор.

OTP: стандартная библиотека архитектуры

OTP (Open Telecom Platform) — набор проверенных абстракций поверх процессов. Главные — GenServer (сервер с состоянием) и Supervisor (наблюдатель, перезапускающий детей). Плюс удобные Agent, Task, Registry, DynamicSupervisor.

GenServer — сервер с состоянием

GenServer инкапсулирует тот самый «цикл с состоянием», добавляя стандартный протокол вызовов, обработку системных сообщений, отладку и интеграцию с супервизорами. Вы реализуете колбэки, OTP — всю рутину.

defmodule KV do
  use GenServer

  ## --- Клиентский API (публичный интерфейс, работает в процессе вызывающего) ---

  def start_link(opts) do
    GenServer.start_link(__MODULE__, %{}, name: Keyword.get(opts, :name, __MODULE__))
  end

  def put(server \\ __MODULE__, key, value), do: GenServer.cast(server, {:put, key, value})
  def get(server \\ __MODULE__, key), do: GenServer.call(server, {:get, key})

  ## --- Серверные колбэки (работают ВНУТРИ процесса GenServer) ---

  @impl true
  def init(state), do: {:ok, state}          # начальное состояние

  @impl true
  def handle_call({:get, key}, _from, state) do
    # call — синхронный: клиент ждёт ответа. {:reply, ЧТО_ВЕРНУТЬ, НОВОЕ_СОСТОЯНИЕ}
    {:reply, Map.get(state, key), state}
  end

  @impl true
  def handle_cast({:put, key, value}, state) do
    # cast — асинхронный: клиент не ждёт. {:noreply, НОВОЕ_СОСТОЯНИЕ}
    {:noreply, Map.put(state, key, value)}
  end
end

call vs cast — важнейшее различие:

  • callсинхронный: клиент блокируется до ответа (по умолчанию таймаут 5 сек). Используйте, когда нужен результат или обратное давление (клиент не забежит вперёд сервера).
  • castасинхронный: «выстрелил и забыл», ответа нет. Быстрее, но без гарантий обработки и без естественного backpressure. Злоупотребление cast — частая причина переполнения почтового ящика под нагрузкой.

Правило: предпочитайте call, если нет веской причины для cast. call даёт обратное давление и заметность ошибок.

Другие полезные колбэки: handle_info/2 (произвольные сообщения, не через call/cast — например :DOWN от мониторов, таймеры), terminate/2 (уборка при завершении), handle_continue/2 (продолжение инициализации после init, чтобы не блокировать старт).

Supervisor — тот, кто перезапускает

Супервизор запускает дочерние процессы, следит за ними через link+trap_exit и перезапускает при падении по заданной стратегии. Прикладная логика и логика восстановления разделены.

defmodule MyApp.Application do
  use Application

  @impl true
  def start(_type, _args) do
    children = [
      {KV, name: KV},                          # запустить наш GenServer
      {Registry, keys: :unique, name: MyApp.Registry},
      {DynamicSupervisor, name: MyApp.WorkerSup, strategy: :one_for_one}
    ]

    Supervisor.start_link(children, strategy: :one_for_one, name: MyApp.Supervisor)
  end
end

Стратегии перезапуска:

  • :one_for_one — упал один ребёнок → перезапускается только он. Самая частая. Используйте, когда дети независимы.
  • :one_for_all — упал один → перезапускаются все дети. Когда дети сильно связаны и должны стартовать как единое целое.
  • :rest_for_one — упал ребёнок → перезапускается он и все, кто объявлен после него. Когда есть зависимость по порядку (B зависит от A).

Restart-политика ребёнка (restart:): :permanent (перезапускать всегда — для долгоживущих сервисов, по умолчанию), :temporary (никогда не перезапускать — задача выполнилась и всё), :transient (перезапускать только при ненормальном завершении, не при :normal).

Защита от бесконечных рестартовmax_restarts (по умолчанию 3) за max_seconds (по умолчанию 5). Если ребёнок падает чаще — супервизор сдаётся и падает сам, эскалируя проблему вверх по дереву. Это осознанный дизайн: если что-то падает без конца, чинить надо не здесь, а уровнем выше.

Деревья супервизии

Реальные приложения — это деревья: супервизоры супервизоров. Отказ изолируется на минимально возможном уровне, а если и он не справляется — эскалирует наверх, вплоть до перезапуска целого поддерева.

Проектирование дерева — это проектирование отказоустойчивости: вы решаете, что с чем падает и в каком порядке восстанавливается.

Удобные абстракции

Agent — простое разделяемое состояние

Когда нужно просто «хранилище состояния» без своей логики обработки сообщений — Agent короче, чем полный GenServer:

{:ok, agent} = Agent.start_link(fn -> %{} end)
Agent.update(agent, &Map.put(&1, :count, 1))
Agent.get(agent, & &1.count)     # => 1

Под капотом — тот же GenServer. Agent хорош для простого состояния; как только появляется нетривиальная логика — переходите на GenServer.

Task — асинхронные вычисления и параллелизм

Task — для «сделать что-то конкурентно и, возможно, забрать результат»:

# Запустить и дождаться результата
task = Task.async(fn -> heavy_computation() end)
# ... параллельно делаем что-то ещё ...
result = Task.await(task, 10_000)

# Распараллелить обработку коллекции (например, N HTTP-запросов сразу)
[1, 2, 3, 4]
|> Task.async_stream(&fetch_url/1, max_concurrency: 4, timeout: 30_000)
|> Enum.map(fn {:ok, result} -> result end)

Task.async_stream/3 — рабочая лошадка параллельной обработки: контролируемая конкурентность (max_concurrency), backpressure (не запустит больше, чем нужно), таймауты. Для «выстрелил и забыл» под супервизией — Task.Supervisor.

Registry — именование и поиск процессов

Как найти процесс, PID которого заранее неизвестен (например, «процесс комнаты чата №42»)? Registry — это встроенный key→pid реестр:

# в дереве супервизии: {Registry, keys: :unique, name: MyApp.Registry}

# процесс регистрирует себя под ключом при старте:
def start_link(room_id) do
  GenServer.start_link(__MODULE__, room_id,
    name: {:via, Registry, {MyApp.Registry, {:room, room_id}}})
end

# найти процесс по ключу:
[{pid, _}] = Registry.lookup(MyApp.Registry, {:room, 42})

keys: :unique — один процесс на ключ; keys: :duplicate — много (удобно для pub/sub: подписчики регистрируются под темой). Registry сам чистит записи упавших процессов через мониторы.

DynamicSupervisor — процессы по требованию

Обычный Supervisor знает детей заранее. Когда процессы создаются динамически (комната чата на каждый вход, воркер на каждую задачу) — нужен DynamicSupervisor:

# в дереве: {DynamicSupervisor, name: MyApp.RoomSup, strategy: :one_for_one}

def open_room(room_id) do
  spec = {RoomServer, room_id}
  DynamicSupervisor.start_child(MyApp.RoomSup, spec)
end

Комбо DynamicSupervisor + Registry — классический рецепт: динамически создавать именованные процессы и находить их по бизнес-ключу. Так строят чаты, игровые сессии, per-user воркеры.

Конвейеры данных: GenStage, Flow, Broadway

Когда данные текут потоком (события, сообщения из очереди, строки лога), а производитель и потребитель работают с разной скоростью, ключевая проблема — обратное давление (backpressure): не дать быстрому производителю завалить медленного потребителя и выесть память.

GenStage решает это через спрос (demand): потребитель запрашивает у производителя ровно столько событий, сколько готов обработать. Производитель не шлёт больше запрошенного. Backpressure встроен в протокол.

Стадии бывают: producer, consumer, producer_consumer (и то и другое — середина конвейера).

Flow — поверх GenStage — даёт «параллельный Enum»: map/filter/reduce над большим потоком с автоматическим распараллеливанием по ядрам и партиционированием.

File.stream!("huge.csv")
|> Flow.from_enumerable()
|> Flow.map(&parse_line/1)
|> Flow.partition()                 # перераспределить по ключу между стадиями
|> Flow.reduce(fn -> %{} end, &aggregate/2)
|> Enum.to_list()

Broadway — высокоуровневый фреймворк для обработки сообщений из очередей (Amazon SQS, RabbitMQ, Kafka, Google Pub/Sub) с готовыми фичами: backpressure, батчинг, автоматическое подтверждение (ack), устойчивость к сбоям, ограничение скорости. Это то, что берут в прод для потоковой обработки.

defmodule MyApp.Pipeline do
  use Broadway

  def start_link(_opts) do
    Broadway.start_link(__MODULE__,
      name: __MODULE__,
      producer: [module: {BroadwaySQS.Producer, queue_url: "..."}, concurrency: 2],
      processors: [default: [concurrency: 10]],          # 10 конкурентных обработчиков
      batchers: [default: [batch_size: 100, concurrency: 2]]  # батчами по 100
    )
  end

  @impl true
  def handle_message(_processor, message, _context) do
    # обработать одно сообщение; при исключении Broadway сам не подтвердит его
    Message.update_data(message, &process/1)
  end

  @impl true
  def handle_batch(_batcher, messages, _batch_info, _context) do
    # обработать батч (например, одна пакетная запись в БД)
    messages
  end
end

Когда что: нужен полный контроль над стадиями — GenStage; надо «параллельно посчитать над большим набором» — Flow; читаете из брокера очередей в проде — Broadway (почти всегда правильный выбор).

Распределённость: несколько нод

BEAM изначально распределённая. Несколько экземпляров («ноды») соединяются в кластер, и процессы на разных машинах общаются тем же send/receive — прозрачно.

# запуск именованной ноды
# iex --sname node1 --cookie secret
# iex --sname node2 --cookie secret

Node.connect(:"node1@host")     # соединить ноды
Node.list()                     # список подключённых нод

# отправить сообщение процессу на другой ноде — тот же send
send({ProcessName, :"node2@host"}, {:hello, self()})

# выполнить функцию на удалённой ноде
:erpc.call(:"node2@host", Module, :function, [args])

Поверх строят кластеризацию (libcluster для авто-обнаружения нод в k8s/облаке), распределённый pub/sub (Phoenix.PubSub), реплицируемое состояние (Phoenix.Tracker для présence), горизонтальное распределение процессов (Horde).

Осторожно: распределённость не отменяет CAP-теорему. Кластер BEAM по умолчанию предполагает надёжную сеть (полносвязный меш) и не решает за вас консистентность — сетевые разделы (netsplit) реальны. Для критичных к согласованности данных полагайтесь на БД, а не на распределённое состояние в памяти.

Горячее обновление кода (hot code reloading)

BEAM умеет заменять код работающего приложения без остановки — историческая суперспособность из телекома (обновить станцию, не роняя звонки). Механизм существует (Code, релизы с appup), но на практике в вебе его применяют редко: проще выкатить новый инстанс и переключить трафик. Знать о существовании стоит; для типичного сервиса это не первоочередной инструмент.

Наблюдаемость процессов

  • :observer.start() — графический инспектор: дерево супервизии, процессы, память, загрузка планировщиков (в IEx на dev-машине).
  • recon — безопасные для прода инструменты диагностики: найти процессы с раздутым mailbox, утечки, «горячие» функции.
  • Process.info(pid, :message_queue_len) — быстрый способ проверить, не растёт ли очередь сообщений (главный симптом того, что GenServer не успевает).
  • :telemetry — стандартные события для метрик (подробно в главе про наблюдаемость).

Главный симптом проблемы конкурентности в Elixir — растущая очередь сообщений у GenServer: значит, он узкое место. Лечится вынесением работы в отдельные процессы/Task, переходом с cast на call (для backpressure) или шардированием состояния.

Типичные ловушки

  • GenServer как узкое место. Один процесс обрабатывает сообщения последовательно. Если в него «стучатся» все — он станет бутылочным горлышком. Дробите состояние (шардинг по ключу через Registry) или выносите тяжёлую работу.
  • Тяжёлая работа в handle_call. Пока колбэк считает, сервер не отвечает другим. Долгую работу делайте в Task/отдельном процессе.
  • Злоупотребление cast. Без backpressure почтовый ящик растёт под нагрузкой до OOM. По умолчанию — call.
  • Большое состояние в процессе. Гигантский map в GenServer → дорогой GC, дорогие сообщения. Держите состояние компактным; большие данные — в ETS/БД.
  • Блокирующая инициализация init/1. Долгий init тормозит старт всего дерева супервизии. Используйте handle_continue.
  • Атомы из внешнего ввода (String.to_atom) — не собираются GC, утечка/DoS.

Итог: как мыслить

  1. Разбейте задачу на изолированные единицы состояния — каждая станет процессом.
  2. Спроектируйте дерево супервизии: что перезапускается, вместе или порознь, в каком порядке.
  3. Ожидаемые ошибки — значениями ({:error, _}); неожиданные — пусть падает, супервизор поднимет.
  4. Для потоков данных — backpressure через GenStage/Broadway, а не «cast в бесконечность».
  5. Наблюдайте за длиной очередей — это пульс вашей конкурентной системы.

Источники

  • Elixir — Processes, GenServer, Supervisor
  • Книга «Elixir in Action» (Saša Jurić) — лучшее объяснение OTP и «let it crash».
  • Broadway, GenStage, Flow
  • Доклады Saša Jurić «The Soul of Erlang and Elixir» — интуиция про надёжность BEAM.

Что дальше

Тестирование — ExUnit, async-тесты, doctests, моки и property-based.

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов