Исходный код Распределённые задачи и метки
В этой главе мы вернёмся к приложению :kv и добавим уровень маршрутизации, который позволит нам распределять запросы между узлами на основе имени корзины.
Уровень маршрутизации будет получать таблицу маршрутизации следующего формата:
[
{?a..?m, :"foo@computer-name"},
{?n..?z, :"bar@computer-name"}
]
Маршрутизатор проверит первый байт имени корзины по таблице и перенаправит запрос на соответствующий узел на основе этого байта. Например, корзина, начинающаяся с буквы «а» (?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, которые мы могли бы использовать в нашей реализации:
Мы могли бы использовать модуль Erlang :erpc для выполнения функций на удалённом узле. Внутри сессии
bar@computer-nameвыше вы можете вызвать:erpc.call(:"foo@computer-name", Hello, :world, []), и она напечатает «hello world»Мы могли бы запустить сервер на другом узле и отправлять запросы на этот узел через API
GenServer. Например, вы можете вызвать сервер на удалённом узле, используяGenServer.call({name, node}, arg)или передав PID удалённого процесса в качестве первого аргумента.Мы могли бы использовать задачи, о которых мы узнали в предыдущей главе, так как они могут запускаться как на локальных, так и на удалённых узлах.
У перечисленных вариантов есть разные свойства. 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
- Модуль global Erlang, который может предоставлять глобальные имена и глобальные блокировки, обеспечивая уникальные имена и уникальные блокировки во всем кластере машин
- Модуль pg Erlang, который позволяет процессам присоединяться к различным группам, общим для всего кластера
- Проект Phoenix PubSub, который предоставляет распределённую систему обмена сообщениями и распределённую систему присутствия для отслеживания пользователей и процессов в кластере
Вы также найдете много библиотек для создания распределённых систем в рамках экосистемы Erlang. Сейчас пора вернуться к нашему простому распределённому хранилищу ключей-значений и узнать, как его настроить и упаковать для производства.
© 2012-2024 The Elixir Team
Licensed under the Apache License, Version 2.0.
https://hexdocs.pm/elixir/1.17.2/distributed-tasks.html