Spec-Zone.ru › Elixir 1.16

Источник Распределённые задачи и теги

В этой главе мы вернёмся к применению :kv и добавим слой маршрутизации, который позволит нам распределять запросы между узлами на основе имени корзины.

Слой маршрутизации будет получать таблицу маршрутизации следующего формата:

[
  {?a..?m, :"foo@computer-name"},
  {?n..?z, :"bar@computer-name"}
]

Маршрутизатор будет проверять первый байт имени корзины по таблице и перенаправлять на соответствующий узел на основе этого. Например, корзина, начинающаяся с буквы "a" (?a представляет собой код Юникода буквы "a"), будет перенаправлена на узел foo@computer-name.

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

Примечание: в этой главе мы будем использовать два узла на одной машине. Вы можете использовать две (или более) разные машины в одной сети, но вам необходимо выполнить некоторые подготовительные действия. Во-первых, вам необходимо убедиться, что на всех машинах есть файл ~/.erlang.cookie с точно таким же значением. Затем необходимо гарантировать, что epmd запущен на порту, который не заблокирован (можно запустить epmd -d для отладки).

Наш первый распределённый код

Elixir поставляется с возможностями подключения узлов и обмена информацией между ними. Фактически, мы используем те же концепции процессов, обмена сообщениями и получения сообщений при работе в распределённой среде, потому что процессы Elixir являются *прозрачными по местоположению*. Это означает, что при отправке сообщения не имеет значения, находится ли получающий процесс на том же узле или на другом узле, виртуальная машина сможет доставить сообщение в обоих случаях.

Для запуска распределённого кода нам необходимо запустить виртуальную машину с именем. Имя может быть коротким (при работе в одной сети) или длинным (требуется полный адрес компьютера). Давайте запустим новую сессию IEx:

$ iex --sname foo

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

Interactive Elixir - press Ctrl+C to exit (type h() ENTER for help)
iex(foo@jv)1>

Мой компьютер называется jv, поэтому в примере выше я вижу foo@jv, но у вас будет другой результат. В следующих примерах мы будем использовать foo@computer-name, и вы должны обновить их соответствующим образом, когда будете пробовать код.

Давайте определим модуль с именем Hello в этой оболочке:

iex> defmodule Hello do
...>   def world, do: IO.puts "hello world"
...> end

Если у вас есть другой компьютер в той же сети с установленным Erlang и Elixir, вы можете запустить другую оболочку на нём. Если нет, вы можете запустить другую сессию IEx в другом терминале. В любом случае, присвойте ему короткое имя bar:

$ iex --sname bar

Обратите внимание, что внутри этой новой сессии IEx мы не можем получить доступ к Hello.world/0:

iex> Hello.world
** (UndefinedFunctionError) function Hello.world/0 is undefined (module Hello is not available)
    Hello.world()

Однако мы можем запустить новый процесс на foo@computer-name из bar@computer-name! Давайте попробуем (где @computer-name - то, что вы видите локально):

iex> Node.spawn_link(:"foo@computer-name", fn -> Hello.world() end)
#PID<9014.59.0>
hello world

Elixir запустил процесс на другом узле и вернул его PID. Код затем выполнился на другом узле, где функция Hello.world/0 существует, и вызвал эту функцию. Обратите внимание, что результат «hello world» был напечатан на текущем узле bar, а не на foo. Другими словами, сообщение, которое должно быть напечатано, было отправлено обратно с foo на bar. Это происходит потому, что процесс, запущенный на другом узле (foo), знает, что весь вывод должен быть отправлен обратно на исходный узел!

Мы можем отправлять и получать сообщения из PID, возвращённого Node.spawn_link/2, как обычно. Давайте попробуем быстрый пример ping-pong:

iex> pid = Node.spawn_link(:"foo@computer-name", fn ->
...>   receive do
...>     {:ping, client} -> send(client, :pong)
...>   end
...> end)
#PID<9014.59.0>
iex> send(pid, {:ping, self()})
{:ping, #PID<0.73.0>}
iex> flush()
:pong
:ok

Из нашего быстрого исследования мы можем сделать вывод, что мы должны использовать Node.spawn_link/2 для запуска процессов на удалённом узле каждый раз, когда нам нужно выполнить распределённое вычисление. Однако, на протяжении всего этого руководства мы узнали, что запускать процессы за пределами деревьев управления следует избегать, если это возможно, поэтому нам нужно искать другие варианты.

Существует три лучших альтернативы Node.spawn_link/2, которые мы могли бы использовать в нашей реализации:

  1. Мы могли бы использовать модуль Erlang :erpc для выполнения функций на удалённом узле. Внутри оболочки bar@computer-name выше вы можете вызвать :erpc.call(:"foo@computer-name", Hello, :world, []), и она напечатает «hello world»

  2. Мы могли бы запустить сервер на другом узле и отправлять запросы на этот узел через API GenServer. Например, вы можете вызвать сервер на удалённом узле, используя GenServer.call({name, node}, arg) или передавая PID удалённого процесса в качестве первого аргумента

  3. Мы могли бы использовать задачи, о которых мы узнали в предыдущей главе, так как они могут запускаться как на локальных, так и на удалённых узлах

У перечисленных вариантов есть разные свойства. GenServer будет сериализовать ваши запросы на одном сервере, а задачи фактически выполняются асинхронно на удалённом узле, единственной точкой сериализации является запуск, выполняемый диспетчером.

Для нашего слоя маршрутизации мы будем использовать задачи, но вы можете изучить и другие альтернативы.

async/await

До сих пор мы изучали задачи, которые запускаются и выполняются изолированно, не обращая внимания на их возвращаемое значение. Однако иногда полезно запустить задачу для вычисления значения и прочитать её результат позже. Для этого задачи также предоставляют шаблон async/await:

task = Task.async(fn -> compute_something_expensive() end)
res  = compute_something_else()
res + Task.await(task)

async/await обеспечивает очень простой механизм для одновременного вычисления значений. Кроме того, async/await также может использоваться с тем же Task.Supervisor, что мы использовали в предыдущих главах. Нам просто нужно вызвать Task.Supervisor.async/2 вместо Task.Supervisor.start_child/2 и использовать Task.await/2, чтобы прочитать результат позже.

Распределённые задачи

Распределённые задачи точно такие же, как управляемые задачи. Единственное отличие заключается в том, что мы передаём имя узла при запуске задачи в диспетчере. Откройте lib/kv/supervisor.ex из приложения :kv. Давайте добавим диспетчера задач как последнего дочернего элемента дерева:

{Task.Supervisor, name: KV.RouterTasks},

Теперь давайте запустим снова два узла с именами, но внутри приложения :kv:

$ iex --sname foo -S mix
$ iex --sname bar -S mix

Изнутри bar@computer-name, мы теперь можем запустить задачу напрямую на другом узле через диспетчер:

iex> task = Task.Supervisor.async({KV.RouterTasks, :"foo@computer-name"}, fn ->
...>   {:ok, node()}
...> end)
%Task{
  mfa: {:erlang, :apply, 2},
  owner: #PID<0.122.0>,
  pid: #PID<12467.88.0>,
  ref: #Reference<0.0.0.400>
}
iex> Task.await(task)
{:ok, :"foo@computer-name"}

Наша первая распределённая задача извлекает имя узла, на котором выполняется задача. Обратите внимание, что мы передали анонимную функцию Task.Supervisor.async/2, но в распределённых случаях предпочтительнее явно указать модуль, функцию и аргументы:

iex> task = Task.Supervisor.async({KV.RouterTasks, :"foo@computer-name"}, Kernel, :node, [])
%Task{
  mfa: {Kernel, :node, 0},
  owner: #PID<0.122.0>,
  pid: #PID<12467.89.0>,
  ref: #Reference<0.0.0.404>
}
iex> Task.await(task)
:"foo@computer-name"

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

С этими знаниями давайте наконец напишем код маршрутизации.

Слой маршрутизации

Создайте файл в lib/kv/router.ex со следующим содержимым:

defmodule KV.Router do
  @doc """
  Dispatch the given `mod`, `fun`, `args` request
  to the appropriate node based on the `bucket`.
  """
  def route(bucket, mod, fun, args) do
    # Get the first byte of the binary
    first = :binary.first(bucket)

    # Try to find an entry in the table() or raise
    entry =
      Enum.find(table(), fn {enum, _node} ->
        first in enum
      end) || no_entry_error(bucket)

    # If the entry node is the current node
    if elem(entry, 1) == node() do
      apply(mod, fun, args)
    else
      {KV.RouterTasks, elem(entry, 1)}
      |> Task.Supervisor.async(KV.Router, :route, [bucket, mod, fun, args])
      |> Task.await()
    end
  end

  defp no_entry_error(bucket) do
    raise "could not find entry for #{inspect bucket} in table #{inspect table()}"
  end

  @doc """
  The routing table.
  """
  def table do
    # Replace computer-name with your local machine name
    [{?a..?m, :"foo@computer-name"}, {?n..?z, :"bar@computer-name"}]
  end
end

Давайте напишем тест, чтобы проверить, что наш маршрутизатор работает. Создайте файл с именем test/kv/router_test.exs со следующим содержимым:

defmodule KV.RouterTest do
  use ExUnit.Case, async: true

  test "route requests across nodes" do
    assert KV.Router.route("hello", Kernel, :node, []) ==
             :"foo@computer-name"
    assert KV.Router.route("world", Kernel, :node, []) ==
             :"bar@computer-name"
  end

  test "raises on unknown entries" do
    assert_raise RuntimeError, ~r/could not find entry/, fn ->
      KV.Router.route(<<0>>, Kernel, :node, [])
    end
  end
end

Первый тест вызывает Kernel.node/0, которая возвращает имя текущего узла, основанные на именах корзин "hello" и "world". Согласно нашей таблице маршрутизации, мы должны получить foo@computer-name и bar@computer-name в ответ соответственно.

Второй тест проверяет, что код выбросит исключение для неизвестных записей.

Для выполнения первого теста необходимо запустить два узла. Перейдите в apps/kv и перезапустите узел с именем bar, который будет использоваться тестами.

$ iex --sname bar -S mix

И теперь запустите тесты с помощью:

$ elixir --sname foo -S mix test

Тест должен пройти.

Фильтры и теги тестов

Хотя наши тесты проходят, наша структура тестирования становится более сложной. В частности, выполнение тестов только с помощью mix test приводит к сбоям в нашем наборе тестов, так как наш тест требует подключения к другому узлу.

К счастью, ExUnit поставляется со средствами для тегов тестов, что позволяет нам запускать определённые обратные вызовы или даже фильтровать тесты на основе этих тегов. В предыдущей главе мы уже использовали тег :capture_log, семантика которого задана самим ExUnit.

На этот раз давайте добавим тег :distributed к test/kv/router_test.exs:

@tag :distributed
test "route requests across nodes" do

Запись @tag :distributed эквивалентна записи @tag distributed: true.

После того, как тест был правильно помечен тегом, мы можем проверить, доступен ли узел в сети, и если нет, мы можем исключить все распределённые тесты. Откройте test/test_helper.exs внутри приложения :kv и добавьте следующее:

exclude =
  if Node.alive?(), do: [], else: [distributed: true]

ExUnit.start(exclude: exclude)

Теперь запустите тесты с помощью mix test:

$ mix test
Excluding tags: [distributed: true]

.......

Finished in 0.05 seconds
9 tests, 0 failures, 1 excluded

На этот раз все тесты прошли, и ExUnit предупредил нас о том, что распределённые тесты были исключены. Если вы запустите тесты с помощью $ elixir --sname foo -S mix test, один дополнительный тест должен быть запущен и успешно пройден, при условии, что узел bar@computer-name доступен.

Команда mix test также позволяет нам динамически включать и исключать теги. Например, мы можем запустить $ mix test --include distributed для запуска распределённых тестов независимо от значения, установленного в test/test_helper.exs. Мы также можем передать --exclude для исключения определённого тега из командной строки. Наконец, --only может использоваться для запуска только тестов с определённым тегом:

$ elixir --sname foo -S mix test --only distributed

Дополнительную информацию о фильтрах, тегах и значениях по умолчанию для тегов можно найти в документации модуля ExUnit.Case.

Объединение всего вместе

Теперь, когда наш механизм маршрутизации настроен, измените KVServer для использования маршрутизатора. Замените функцию lookup/2 в KVServer.Command следующим:

defp lookup(bucket, callback) do
  case KV.Registry.lookup(KV.Registry, bucket) do
    {:ok, pid} -> callback.(pid)
    :error -> {:error, :not_found}
  end
end

на это:

defp lookup(bucket, callback) do
  case KV.Router.route(bucket, KV.Registry, :lookup, [KV.Registry, bucket]) do
    {:ok, pid} -> callback.(pid)
    :error -> {:error, :not_found}
  end
end

Вместо непосредственного поиска в реестре мы используем маршрутизатор для сопоставления конкретного узла. Затем мы получаем pid, который может принадлежать любому процессу в нашем кластере. Отныне запросы GET, PUT и DELETE направляются на соответствующий узел.

Давайте также убедимся, что при создании нового ведра оно оказывается на правильном узле. Замените функцию run/1 в KVServer.Command, соответствующую команде :create, на следующее:

def run({:create, bucket}) do
  case KV.Router.route(bucket, KV.Registry, :create, [KV.Registry, bucket]) do
    pid when is_pid(pid) -> {:ok, "OK\r\n"}
    _ -> {:error, "FAILED TO CREATE BUCKET"}
  end
end

Теперь, если вы запустите тесты, вы увидите, что существующий тест, проверяющий взаимодействие с сервером, завершится сбоем, поскольку он попытается использовать таблицу маршрутизации. Чтобы исправить этот сбой, измените test_helper.exs для приложения :kv_server, как мы сделали для :kv, и добавьте @tag :distributed к этому тесту тоже:

@tag :distributed
test "server interaction", %{socket: socket} do

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

Подведение итогов

Мы лишь коснулись поверхности того, что возможно при распределённой обработке.

Во всех наших примерах мы полагались на способность Erlang автоматически подключаться к узлам всякий раз, когда поступает запрос. Например, когда мы вызываем Node.spawn_link(:"foo@computer-name", fn -> Hello.world() end), Erlang автоматически подключается к указанному узлу и запускает новый процесс. Однако вы также можете использовать более явный подход к подключениям, используя Node.connect/1 и Node.disconnect/1.

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

В производственной среде узлы могут подключаться и отключаться в любое время. В таких сценариях необходимо обеспечить обнаружение узлов. Библиотеки, такие как libcluster и dns_cluster, предоставляют несколько стратегий обнаружения узлов с использованием DNS, Kubernetes и т.д.

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

Эти темы могут показаться сложными на первый взгляд, но помните, что большинство фреймворков Elixir абстрагируют эти проблемы от вас. Например, при использовании фреймворка Phoenix его абстракции plug-and-play обрабатывают отправку сообщений и отслеживание того, как пользователи присоединяются к и покидают кластер. Однако, если вас интересуют распределённые системы, есть много чего изучить. Вот дополнительные ссылки:

  • Отличная глава Distribunomicon из Learn You Some Erlang
  • Модуль Erlang global, который может предоставлять глобальные имена и глобальные блокировки, позволяя использовать уникальные имена и блокировки в целом кластере машин
  • Модуль Erlang pg, который позволяет процессам присоединяться к различным группам, общим для всего кластера
  • Проект Phoenix PubSub, который предоставляет распределённую систему обмена сообщениями и распределённую систему присутствия для отслеживания пользователей и процессов в кластере

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

← Предыдущая страница Doctests, шаблоны и с помощью with
Следующая страница → Настройка и релизы

Скачать версию ePub

Создано с помощью ExDoc (v0.32.2) для Elixir programming language

© 2012-2024 The Elixir Team
Licensed under the Apache License, Version 2.0.
https://hexdocs.pm/elixir/1.16.3/distributed-tasks.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API