NameNode HA: QJM quorum, Standby NameNode и автоматический failover
Глубокий разбор архитектуры High Availability HDFS: проблема SPOF, синхронизация через QJM и JournalNodes, автоматический failover через ZKFC и ZooKeeper, механизмы fencing, поведение DataNode, конфигурация Spark-клиента и диагностика failover-событий в Spark UI.
1. Архитектурная проблема Single Point of Failure в классическом HDFS¶
Прежде чем понять зачем нужен HA, необходимо разобраться в том, что именно происходит при отказе единственной NameNode в классической Hadoop-конфигурации и почему это событие так катастрофично для Spark-кластера.
Анатомия NameNode: что она хранит и почему это важно¶
NameNode - это мозг HDFS. Всё дерево директорий и файлов файловой системы, метаданные каждого файла и полная картина физического расположения блоков на DataNode - всё это хранится в оперативной памяти NameNode. Это сознательное архитектурное решение: доступ к RAM принципиально быстрее обращения к диску, и именно поэтому HDFS может обслуживать тысячи параллельных запросов на получение блочных адресов с минимальной задержкой.
Конкретно в памяти NameNode живут три ключевых структуры:
FSImage - снапшот дерева директорий и метаданных файлов на конкретный момент времени. Это большой бинарный файл, который периодически записывается на диск как чекпоинт. При запуске NameNode загружает именно его - это фундамент, с которого начинается восстановление.
EditLog - журнал всех изменений файловой системы после последнего FSImage. Каждая операция (создание файла, переименование, удаление) записывается сюда прежде чем применяется к in-memory структурам. Если NameNode упадёт и поднимется, она воспроизведёт EditLog поверх FSImage и восстановит актуальное состояние.
Block Map - маппинг каждого блока файла на список DataNode, хранящих его реплики. Block Map не персистируется - он восстанавливается из heartbeat-отчётов DataNode при каждом старте NameNode. Именно поэтому запуск NameNode после аварии занимает столько времени: нужно дождаться Block Report от всех DataNode.
При типичной нагрузке Data Lake с несколькими сотнями миллионов файлов NameNode потребляет 50–100 GB оперативной памяти только на хранение метаданных. Это дорогостоящий ресурс, но альтернативы нет: если метаданные перенести на диск, задержка каждого getBlockLocations() вырастет с микросекунд до миллисекунд, и весь кластер деградирует.
Что происходит при падении NameNode¶
Когда NameNode становится недоступен - неважно, по какой причине (плановая перезагрузка, OOM, аппаратный отказ, долгий GC) - для Spark-кластера это означает следующее:
Все новые операции чтения мгновенно блокируются. Spark Driver перед планированием задач вызывает listStatus() и getBlockLocations() на NameNode. Если NameNode не отвечает, Driver не может получить физические адреса блоков и не может создать InputSplit'ы для Executor'ов. Сессия зависает на ожидании ответа.
Все новые операции записи невозможны. Попытка создать новый файл или директорию - create(), mkdirs() - немедленно завершается ошибкой StandbyException или таймаутом соединения.
Уже запущенные задачи прерываются. Executor'ы в процессе чтения блоков HDFS время от времени обращаются к NameNode для переаллокации блоков при ошибках DataNode или для обновления lease. Если NameNode недоступен достаточно долго - активные задачи тоже начинают падать.
Checkpoint Secondary NameNode не помогает. В классической конфигурации Secondary NameNode занимается только периодическим слиянием FSImage с EditLog - она не является горячим резервом и не может немедленно принять трафик. Восстановление из Secondary NameNode занимает минуты или часы в зависимости от размера EditLog.
Схема демонстрирует ключевой парадокс классического HDFS: физические данные на DataNode остаются нетронутыми и доступными, но кластер фактически мёртв, потому что нет узла, способного ответить на вопрос «где лежит этот файл». Данные есть, но добраться до них нельзя.
Метрики отказа: RTO и RPO в банковском и телеком секторе¶
Для понимания масштаба проблемы рассмотрим конкретные числа. При классической конфигурации HDFS с одним NameNode:
- Время обнаружения отказа: 30–120 секунд (YARN/monitoring системы)
- Время восстановления NameNode из чекпоинта: 5–30 минут (зависит от размера EditLog и объёма Block Map)
- Итоговый RTO (Recovery Time Objective): 10–40 минут недоступности кластера
Для банков и телекомов, где пайплайны обработки транзакций или событий клиентов работают в режиме 24×7, даже 10 минут простоя означают нарушение SLA и потенциальные финансовые потери. Именно поэтому HA NameNode стала не роскошью, а стандартом production HDFS-кластеров.
2. Синхронизация состояния через QJM: Quorum Journal Manager¶
HA-конфигурация HDFS решает проблему SPOF через схему Active/Standby NameNode. Но немедленно возникает следующий вопрос: как обеспечить идентичность состояния обоих NameNode в реальном времени? Ведь если при переключении окажется, что Standby NN отстал на несколько минут от Active - возникнет потеря данных (нарушение RPO).
Проблема Split-Brain: почему нельзя просто реплицировать данные¶
Наивный подход - реплицировать EditLog с Active NN на Standby NN через NFS-шару или сетевой диск. Эта схема кажется простой, но создаёт смертельно опасную проблему Split-Brain.
Split-Brain - ситуация, когда оба NameNode одновременно считают себя Active и начинают обрабатывать клиентские запросы. Это может произойти при сетевом партиционировании: Active NN изолирован от части кластера, ZKFC решает поднять Standby, но старый Active при этом ещё жив и продолжает писать данные.
Результат Split-Brain в HDFS: оба NameNode начинают независимо изменять namespace. Файл создаётся в одном NameNode и не существует в другом. Блоки, которые Active NN считает принадлежащими файлу A, второй NN отдаёт файлу B. Метаданные расходятся, и восстановить их консистентность без потери данных невозможно.
NFS-шара не решает Split-Brain, потому что нет механизма, гарантирующего, что только один NameNode в каждый момент времени выполняет запись.
Архитектура QJM: кластер JournalNode как надёжный общий журнал¶
QJM (Quorum Journal Manager) - это специализированная распределённая система, созданная именно для решения задачи синхронизации EditLog между двумя NameNode без риска Split-Brain.
Принцип работы QJM. Каждый JournalNode - это лёгкий JVM-процесс, который принимает сегменты EditLog от Active NameNode и хранит их на локальном диске. Никаких репликаций между JournalNode нет - они работают независимо.
Active NameNode записывает каждое изменение файловой системы в сегмент EditLog и отправляет его на все JournalNode параллельно. Запись считается подтверждённой, когда большинство JournalNode (кворум) ответило ACK. При трёх JournalNode кворум = 2, при пяти = 3.
Standby NameNode непрерывно опрашивает JournalNode, читает новые сегменты и применяет их к своей копии namespace. Таким образом Standby всегда на несколько секунд (или миллисекунд при быстрой сети) отстаёт от Active - но не на минуты.
Epoch ID: защита от Split-Brain в QJM¶
Ключевой механизм безопасности QJM - это Epoch ID (также называемый Writer ID). Каждый раз, когда NameNode становится Active, он получает новый Epoch ID, который больше всех предыдущих.
JournalNode принимает записи только от NameNode с текущим максимальным Epoch ID. Если «старый» Active NN (изолированный от ZooKeeper) попытается записать в JournalNode EditLog с устаревшим Epoch ID - JournalNode откажет.
Сценарий Split-Brain с Epoch ID:
1. nn1 - Active, Epoch=5, пишет в QJM
2. Сетевой сбой изолирует nn1 от ZooKeeper
3. nn2 становится Active, получает Epoch=6
4. nn2 начинает писать в QJM с Epoch=6
5. nn1 восстанавливается, пробует писать с Epoch=5
6. JournalNode: "Epoch 5 < текущего 6, отказываю"
7. nn1 понимает, что есть более новый Active, уходит в Standby
Это математически гарантированная защита от Split-Brain: кворум JournalNode физически не может принять конкурирующие записи с разными Epoch ID, поскольку записи с меньшим Epoch просто отклоняются на уровне протокола.
Почему нечётное число JournalNode¶
JournalNode должно быть нечётное число: 3, 5 или 7. Это следует из математики кворума:
- 3 JN: допускает отказ 1 узла (кворум = 2 из 3)
- 5 JN: допускает отказ 2 узлов (кворум = 3 из 5)
- 7 JN: допускает отказ 3 узлов (кворум = 4 из 7)
С чётным числом, например 4, ситуация «2 узла доступны, 2 недоступны» не имеет кворума - система блокируется. Нечётное число гарантирует, что большинство всегда определено однозначно.
Standby NameNode как постоянно готовый горячий резерв¶
В HA-конфигурации Standby NN выполняет работу, которая в классической схеме лежала на Secondary NN - периодическое слияние FSImage с EditLog (checkpoint). Это важное отличие: в HA-схеме Secondary NameNode не нужен. Standby NN одновременно поддерживает актуальный namespace и создаёт чекпоинты.
Для создания чекпоинта Standby NN периодически:
- Создаёт новый FSImage из текущего in-memory состояния namespace
- Загружает новый FSImage на Active NN (через HTTP)
- Active NN начинает использовать новый FSImage как базовую точку
Это означает, что EditLog на JournalNode никогда не вырастает до бесконечных размеров - он регулярно «обрезается» после создания чекпоинта.
3. Механизм автоматического Failover: ZKFC и Apache ZooKeeper¶
Синхронизация через QJM гарантирует, что Standby всегда имеет актуальное состояние namespace. Но сам по себе QJM не знает, когда нужно переключиться. Для этого нужен отдельный механизм - автоматический детектор отказов и инициатор переключения.
Apache ZooKeeper как арбитр выбора лидера¶
ZooKeeper - распределённая система координации, обеспечивающая консенсус между узлами. В контексте HDFS HA она выполняет одну специфическую задачу: хранение «эфемерного узла» (ephemeral znode), который сигнализирует о том, какой NameNode является Active.
Эфемерный znode в ZooKeeper существует только пока жива ZooKeeper-сессия клиента, который его создал. Когда клиент (или процесс ZKFC на его стороне) теряет связь с ZooKeeper - znode автоматически удаляется. Это ключевое свойство, на котором строится весь механизм failover.
ZKFailoverController (ZKFC): локальный страж каждого NameNode¶
ZKFC - небольшой Java-демон, запускаемый на каждом хосте с NameNode (обычно как отдельный процесс hdfs zkfc). Он выполняет три функции одновременно:
1. Мониторинг здоровья локального NameNode. ZKFC каждые несколько секунд отправляет RPC-запрос своему NameNode и ждёт ответа. Не любой ответ считается «здоровым» - NameNode должен ответить за разумное время и вернуть статус HEALTHY (не INITIALIZING, не SERVICE_UNHEALTHY).
2. Управление ZooKeeper-сессией. ZKFC поддерживает постоянную сессию с ZooKeeper-кластером. Если NameNode на машине ZKFC1 жив и здоров - ZKFC1 держит сессию активной и не позволяет ephemeral znode /ActiveStandbyElectorLock исчезнуть.
3. Реакция на изменения znode. ZKFC2 (на хосте Standby NN) установил watch на /ActiveStandbyElectorLock. Когда этот znode исчезнет (сессия ZKFC1 прервалась), ZooKeeper немедленно уведомит ZKFC2.
Полный алгоритм автоматического Failover¶
Диаграмма показывает полный путь от момента отказа до восстановления работоспособности кластера. Обратите внимание на ключевой шаг - Fencing - который происходит до transitionToActive(). Это критически важно: нельзя активировать Standby, не убедившись, что старый Active точно не может выполнять операции.
Механизмы Fencing: как гарантировать отключение старого Active¶
Fencing - это процедура принудительного и необратимого отключения предыдущего Active NameNode перед тем, как новый займёт его место. Без fencing возможен Split-Brain: QJM защищает журнал через Epoch ID, но NameNode теоретически мог успеть ответить клиентам на старые запросы.
Hadoop поддерживает несколько методов fencing, настраиваемых через dfs.ha.fencing.methods:
1. sshfence - наиболее распространённый метод. ZKFC выполняет SSH-команду на хосте старого Active NN и принудительно убивает процесс NameNode:
<property>
<name>dfs.ha.fencing.methods</name>
<!-- Порядок важен: сначала пробуем SSH kill, затем no-op если уже убит -->
<value>sshfence(hdfs)
shell(/bin/true)</value>
</property>
<property>
<name>dfs.ha.fencing.ssh.private-key-files</name>
<!-- SSH ключ, которым ZKFC может подключиться к соседней NN -->
<value>/home/hdfs/.ssh/id_rsa</value>
</property>
sshfence пытается подключиться к хосту через SSH и запустить:
fuser -v -k -n tcp <port> # убивает процесс, слушающий на порту NameNode
2. Shell с IPMI/iLO - используется в production для гарантированного отключения. IPMI позволяет физически обесточить сервер или выполнить hard reset:
<property>
<name>dfs.ha.fencing.methods</name>
<value>shell(/usr/local/bin/fence-ipmi.sh ${dfs.ha.fencing.target.host})</value>
</property>
3. Сетевые методы - отключение порта коммутатора через SNMP или API сетевого оборудования. Более надёжный метод, чем SSH, потому что работает даже при зависании ОС хоста.
Правило безопасности для Fencing: если ни один метод fencing не сработал (SSH недоступен, IPMI не ответил), failover не завершается - Standby NN остаётся в Standby. Лучше иметь недоступный кластер, чем допустить Split-Brain и повреждение данных.
Что происходит, если NameNode завис на GC, но не умер¶
Особый сценарий - Java GC Pause в Active NN. JVM может уйти в Full GC на 30–120 секунд. В это время NameNode не отвечает на RPC, ZKFC детектирует отказ и инициирует failover. Но NameNode не мёртв - он просто временно недоступен. После GC он «воскресает» и видит, что его место занято.
Поведение после GC:
- NameNode пытается писать в JournalNode с устаревшим Epoch ID - получает отказ
- NameNode обнаруживает потерю ZK-сессии
- NameNode самостоятельно переходит в Standby: вызывает
transitionToStandby()
Таким образом QJM (через Epoch ID) и ZooKeeper совместно обеспечивают правильное поведение даже при временных зависаниях.
4. Поведение DataNode при переключении мастеров¶
DataNode в HA-кластере ведут себя иначе, чем в классической конфигурации. Это менее очевидная, но важная часть архитектуры.
Двойная отчётность: Heartbeat и Block Report на обе NameNode¶
В классическом HDFS DataNode отправляет Heartbeat и Block Report только на одну NameNode. В HA-конфигурации DataNode знает об обоих NameNode и поддерживает соединения с обоими:
Почему Standby NN должна получать Block Report в реальном времени?
Вспомним: Block Map не персистируется на диске. При перезапуске NameNode её нужно полностью восстановить из Block Report'ов всех DataNode. Это занимает значительное время - при тысячах DataNode и миллиардах блоков это может занять 10–30 минут.
В HA-конфигурации Standby NN уже имеет актуальный Block Map, потому что постоянно получает Heartbeat и Block Report от DataNode. Это означает, что при failover новый Active NN мгновенно знает физическое расположение всех блоков - нет нужды ждать Block Report'ов. Именно здесь кроется разница между RTO в 40 минут (классический HDFS) и 20 секундами (HA HDFS).
Фазы Block Report при инициализации Standby¶
При первом старте Standby NN или после её длительного простоя происходит Initial Block Report:
- DataNode получает новое подключение от Standby NN
- DataNode отправляет полный Block Report (список всех блоков на диске)
- Standby NN строит Block Map из этих отчётов
Инкрементальные Block Report'ы после этого отправляются при каждом изменении: создании нового блока, удалении, репликации. DataNode сохраняет очередь таких инкрементальных отчётов отдельно для Active и Standby NN.
Команды DataNode: только от Active¶
Несмотря на двойную отчётность, команды DataNode (удалить блок, создать реплику, перенести блок) принимаются только от Active NameNode. Standby NN может накапливать очередь команд внутри себя, но исполнение начнётся только после перехода в Active.
Это важно: если Standby NN «знает», что блок нужно удалить (например, файл был удалён через Active NN), но сама не отдаёт эту команду DataNode. После failover новый Active NN выполнит накопленные команды из своей очереди.
5. Конфигурация HDFS Client: прозрачность для Spark¶
Ключевое преимущество HA-архитектуры для разработчиков - полная прозрачность для клиентского кода. Spark-пайплайн, использующий hdfs://nameservice/path/to/data, продолжает работать без изменения кода как до, так и после failover.
hdfs-site.xml: полная конфигурация HA-кластера¶
<!-- /etc/hadoop/conf/hdfs-site.xml - должен быть идентичен на всех узлах кластера -->
<!--
dfs.nameservices: имя кластера HDFS. Клиенты используют это имя
в URI вместо конкретного hostname: hdfs://mycluster/path
Можно иметь несколько nameservice для нескольких независимых кластеров.
-->
<property>
<name>dfs.nameservices</name>
<value>mycluster</value>
</property>
<!--
Присваиваем логические имена двум NameNode: nn1 и nn2.
Эти имена используются только внутри конфигурации, не в DNS.
-->
<property>
<name>dfs.ha.namenodes.mycluster</name>
<value>nn1,nn2</value>
</property>
<!--
dfs.namenode.rpc-address: адрес, по которому DataNode и клиенты
обращаются к NameNode за метаданными.
Порт 8020 - стандартный для HDFS RPC.
-->
<property>
<name>dfs.namenode.rpc-address.mycluster.nn1</name>
<value>namenode1.example.com:8020</value>
</property>
<property>
<name>dfs.namenode.rpc-address.mycluster.nn2</name>
<value>namenode2.example.com:8020</value>
</property>
<!--
dfs.namenode.http-address: веб-интерфейс и HTTP API NameNode.
Порт 9870 (HDFS 3.x). Используется для чекпоинтинга и веб-UI.
-->
<property>
<name>dfs.namenode.http-address.mycluster.nn1</name>
<value>namenode1.example.com:9870</value>
</property>
<property>
<name>dfs.namenode.http-address.mycluster.nn2</name>
<value>namenode2.example.com:9870</value>
</property>
<!--
dfs.namenode.shared.edits.dir: URI кластера JournalNode.
Формат: qjournal://host1:port;host2:port;host3:port/nameservice
Порт 8485 - стандартный JournalNode RPC порт.
Нечётное число JournalNode ОБЯЗАТЕЛЬНО: 3, 5 или 7.
-->
<property>
<name>dfs.namenode.shared.edits.dir</name>
<value>qjournal://jn1.example.com:8485;jn2.example.com:8485;jn3.example.com:8485/mycluster</value>
</property>
<!--
dfs.client.failover.proxy.provider: класс, который реализует логику
выбора Active NameNode на стороне клиента (Spark, HDFS клиент, etc.)
ConfiguredFailoverProxyProvider - стандартная реализация.
Алгоритм:
1. Пробуем соединиться с nn1
2. Если nn1 не отвечает - переходим к nn2
3. При получении StandbyException - переключаемся автоматически
-->
<property>
<name>dfs.client.failover.proxy.provider.mycluster</name>
<value>org.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider</value>
</property>
<!-- Fencing методы -->
<property>
<name>dfs.ha.fencing.methods</name>
<value>sshfence
shell(/bin/true)</value>
</property>
<property>
<name>dfs.ha.fencing.ssh.private-key-files</name>
<value>/home/hdfs/.ssh/id_rsa</value>
</property>
<!-- ZooKeeper для автоматического failover -->
<property>
<name>ha.zookeeper.quorum</name>
<value>zk1.example.com:2181,zk2.example.com:2181,zk3.example.com:2181</value>
</property>
<property>
<name>dfs.ha.automatic-failover.enabled</name>
<value>true</value>
</property>
Тюнинг клиентских таймаутов под Spark¶
По умолчанию HDFS-клиент довольно агрессивен в обнаружении «мёртвой» NameNode и быстром переключении. Для Spark-джобов, которые работают часами, это иногда создаёт ложные срабатывания во время планового обслуживания (rolling restart NameNode).
<!-- Тюнинг таймаутов HDFS HA клиента для длинных Spark-джобов -->
<!--
Количество попыток подключения к NameNode перед переключением на вторую.
Дефолт: 15. Для production с долгими джобами рекомендуется 20-30.
Увеличение делает переключение медленнее, но уменьшает
ложные failover при кратких перегрузках NameNode.
-->
<property>
<name>dfs.client.failover.connection.retries</name>
<value>20</value>
</property>
<!--
dfs.client.failover.max.attempts: максимальное число попыток
подключения к любой NameNode до возврата ошибки клиенту.
Дефолт: 15.
-->
<property>
<name>dfs.client.failover.max.attempts</name>
<value>30</value>
</property>
<!--
Базовая задержка между попытками failover (миллисекунды).
Дефолт: 500 мс. Используется с exponential backoff.
Итоговая задержка: base * 2^(attempt_number)
Попытка 1: 500мс, попытка 2: 1000мс, попытка 3: 2000мс...
-->
<property>
<name>dfs.client.failover.sleep.base.millis</name>
<value>500</value>
</property>
<!--
Максимальная задержка между попытками.
Дефолт: 15 000 мс (15 секунд).
-->
<property>
<name>dfs.client.failover.sleep.max.millis</name>
<value>15000</value>
</property>
<!--
Таймаут RPC-соединения с NameNode.
Дефолт: 20 000 мс. При медленных сетях или нагруженных NN
рекомендуется увеличить до 30-60 секунд.
-->
<property>
<name>ipc.client.connect.timeout</name>
<value>30000</value>
</property>
Как Spark находит Active NameNode¶
Spark при создании SparkSession загружает hdfs-site.xml из classpath и инициализирует ConfiguredFailoverProxyProvider. При первом подключении Spark пробует nn1. Если nn1 отвечает статусом StandbyException (то есть он Standby, а не Active), клиент автоматически переключается на nn2. Результат кешируется: следующий запрос пойдёт напрямую на nn2.
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("hdfs-ha-demo") \
# Spark автоматически подхватывает hdfs-site.xml из $HADOOP_CONF_DIR.
# Если конфиг не в стандартном месте - указываем явно через spark.hadoop.*
.config("spark.hadoop.fs.defaultFS", "hdfs://mycluster") \
.config("spark.hadoop.dfs.nameservices", "mycluster") \
.config("spark.hadoop.dfs.ha.namenodes.mycluster", "nn1,nn2") \
.config("spark.hadoop.dfs.namenode.rpc-address.mycluster.nn1",
"namenode1.example.com:8020") \
.config("spark.hadoop.dfs.namenode.rpc-address.mycluster.nn2",
"namenode2.example.com:8020") \
.config("spark.hadoop.dfs.client.failover.proxy.provider.mycluster",
"org.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider") \
.getOrCreate()
# Путь использует логическое имя nameservice, а не конкретный хост.
# Spark не знает и не должен знать, nn1 или nn2 сейчас Active.
df = spark.read.parquet("hdfs://mycluster/data/events/")
df.show()
Диагностика текущего Active/Standby через административные команды¶
# Проверить статус обоих NameNode
hdfs haadmin -getServiceState nn1
# Вывод: active
hdfs haadmin -getServiceState nn2
# Вывод: standby
# Список всех NameNode с их статусом
hdfs haadmin -getAllServiceState
# Принудительный ручной failover (при плановом обслуживании)
# Делает nn1 Standby, nn2 Active - с прохождением через fencing
hdfs haadmin -failover nn1 nn2
# Проверить состояние ZooKeeper для HA
hdfs zkfc -formatZK # первоначальная инициализация ZK (один раз при настройке)
# Статус через dfsadmin
hdfs dfsadmin -report
6. Влияние Failover на жизненный цикл Spark Job¶
Failover HDFS NameNode - событие, которое Spark-пайплайн в идеале должен переживать незаметно. На практике это зависит от того, в какой фазе находится джоба в момент переключения.
Фазы Spark Job и их уязвимость к Failover¶
Фаза 3 (выполнение задач) - наиболее защищённая фаза. Executor'ы в этот момент читают данные напрямую с DataNode по TCP, минуя NameNode. Failover NameNode проходит для них незаметно - они продолжают читать с тех же DataNode.
Фаза 5 (commit) - наиболее уязвимая фаза. Если failover произойдёт именно в момент выполнения rename в commit фазе - возможно частичное переименование или дублирование. Для защиты следует использовать idempotent committer (FileOutputCommitter v2) и _SUCCESS файл как маркер завершения.
Как читать логи Spark Driver при Failover¶
При failover в логах Spark Driver появляются характерные сообщения. Зная их, можно точно определить когда и насколько долго длился failover:
# Типичный лог Spark Driver во время failover NameNode
# Момент обнаружения проблемы:
24/01/15 10:23:41 WARN ipc.Client: Retrying connect to server:
namenode1.example.com/192.168.1.10:8020. Already tried 0 time(s);
retry policy is RetryUpToMaximumCountWithFixedSleep(maxRetries=1,
sleepTime=1000 MILLISECONDS)
# Клиент переключается на второй NameNode:
24/01/15 10:23:42 INFO ha.RetryInvocationHandler: Exception while
invoking getBlockLocations of class ClientNamenodeProtocolTranslatorPB
over namenode1.example.com/192.168.1.10:8020 after 1 failover attempts.
Trying to failover after sleeping for 0ms.
# Успешное подключение к новому Active:
24/01/15 10:23:43 INFO BlockReaderFactory: Successfully connected to
namenode2.example.com:8020
# Джоба продолжает выполнение:
24/01/15 10:23:43 INFO FileOutputCommitter: Saved output of task
attempt_20240115_0001_m_000042_0 to hdfs://mycluster/output/
Из этих строк видно: 10:23:41 - первая ошибка, 10:23:43 - успешное подключение к nn2. Итого 2 секунды накладных расходов на failover. Для 30-минутной джобы - незаметно. Для стримингового пайплайна с checkpoint каждые 30 секунд - критично.
Анализ Scheduler Delay в Spark UI¶
В Spark UI на вкладке Stages → конкретный Stage → Summary Metrics есть метрика Scheduler Delay. Она показывает время от момента, когда задача была готова к запуску, до момента, когда она действительно началась.
При нормальной работе Scheduler Delay составляет 1–10 мс. Если в момент планирования тасок происходил failover NameNode - Scheduler Delay для тасок этого Stage вырастет до сотен миллисекунд или секунд, пока Driver ждал ответа от NameNode на listStatus(). Это и есть сигнатура failover-события в Spark UI.
Влияние Rolling Restart NameNode на долгие Spark-джобы¶
Плановое обслуживание HDFS кластера требует перезапуска NameNode. При правильно настроенных таймаутах эта процедура должна быть прозрачной для Spark.
Процедура безопасного Rolling Restart NameNode:
# Шаг 1: сначала перезапускаем Standby NN - безопасно,
# Active продолжает обслуживать запросы
# На хосте nn2 (Standby):
sudo systemctl restart hadoop-hdfs-namenode
# Шаг 2: ждём, пока nn2 вернётся в Standby и синхронизируется с QJM
hdfs haadmin -getServiceState nn2
# Ждём: standby
# Шаг 3: плановый failover на nn2
# Переключаем трафик с nn1 на nn2 - прозрачно для клиентов
hdfs haadmin -failover nn1 nn2
# Шаг 4: перезапускаем бывший Active (теперь Standby)
sudo systemctl restart hadoop-hdfs-namenode # на nn1
# Шаг 5: возвращаем Active обратно на nn1 (если нужно)
hdfs haadmin -failover nn2 nn1
7. Практика: симуляция аварии и траблшутинг¶
Docker Compose стенд для изучения HDFS HA¶
Следующий стенд разворачивает минимальный HA-кластер: 2 NameNode, 3 JournalNode, 3-нодовый ZooKeeper ансамбль и 3 DataNode.
# docker-compose.yml - минимальный HDFS HA кластер
services:
# ── ZooKeeper Ensemble (3 нод) ─────────────────────────────────────────
zookeeper1:
image: confluentinc/cp-zookeeper:7.6.0
environment:
ZOOKEEPER_CLIENT_PORT: 2181
ZOOKEEPER_SERVER_ID: 1
ZOOKEEPER_SERVERS: "zookeeper1:2888:3888;zookeeper2:2888:3888;zookeeper3:2888:3888"
ports:
- "2181:2181"
networks:
- hadoop-net
zookeeper2:
image: confluentinc/cp-zookeeper:7.6.0
environment:
ZOOKEEPER_CLIENT_PORT: 2182
ZOOKEEPER_SERVER_ID: 2
ZOOKEEPER_SERVERS: "zookeeper1:2888:3888;zookeeper2:2888:3888;zookeeper3:2888:3888"
networks:
- hadoop-net
zookeeper3:
image: confluentinc/cp-zookeeper:7.6.0
environment:
ZOOKEEPER_CLIENT_PORT: 2183
ZOOKEEPER_SERVER_ID: 3
ZOOKEEPER_SERVERS: "zookeeper1:2888:3888;zookeeper2:2888:3888;zookeeper3:2888:3888"
networks:
- hadoop-net
# ── JournalNode Cluster (3 нод) ────────────────────────────────────────
journalnode1:
image: apache/hadoop:3.3.6
command: ["hdfs", "journalnode"]
volumes:
- ./conf:/etc/hadoop/conf:ro
- jn1-data:/hadoop/dfs/journal
networks:
- hadoop-net
journalnode2:
image: apache/hadoop:3.3.6
command: ["hdfs", "journalnode"]
volumes:
- ./conf:/etc/hadoop/conf:ro
- jn2-data:/hadoop/dfs/journal
networks:
- hadoop-net
journalnode3:
image: apache/hadoop:3.3.6
command: ["hdfs", "journalnode"]
volumes:
- ./conf:/etc/hadoop/conf:ro
- jn3-data:/hadoop/dfs/journal
networks:
- hadoop-net
# ── NameNode 1 (стартует как Active) ───────────────────────────────────
namenode1:
image: apache/hadoop:3.3.6
command:
- /bin/sh
- -c
- |
if [ ! -d /hadoop/dfs/name/current ]; then
hdfs namenode -format -clusterId cid1 -force
hdfs zkfc -formatZK -force
fi
hdfs zkfc &
hdfs namenode
depends_on:
- journalnode1
- journalnode2
- journalnode3
- zookeeper1
ports:
- "9870:9870"
- "8020:8020"
volumes:
- ./conf:/etc/hadoop/conf:ro
- nn1-data:/hadoop/dfs/name
networks:
hadoop-net:
aliases:
- namenode1
# ── NameNode 2 (стартует как Standby) ──────────────────────────────────
namenode2:
image: apache/hadoop:3.3.6
command:
- /bin/sh
- -c
- |
if [ ! -d /hadoop/dfs/name/current ]; then
hdfs namenode -bootstrapStandby -force
fi
hdfs zkfc &
hdfs namenode
depends_on:
- namenode1
ports:
- "9871:9870"
- "8021:8020"
volumes:
- ./conf:/etc/hadoop/conf:ro
- nn2-data:/hadoop/dfs/name
networks:
hadoop-net:
aliases:
- namenode2
datanode1:
image: apache/hadoop:3.3.6
command: ["hdfs", "datanode"]
depends_on: [namenode1]
volumes:
- ./conf:/etc/hadoop/conf:ro
- dn1-data:/hadoop/dfs/data
networks:
- hadoop-net
datanode2:
image: apache/hadoop:3.3.6
command: ["hdfs", "datanode"]
depends_on: [namenode1]
volumes:
- ./conf:/etc/hadoop/conf:ro
- dn2-data:/hadoop/dfs/data
networks:
- hadoop-net
datanode3:
image: apache/hadoop:3.3.6
command: ["hdfs", "datanode"]
depends_on: [namenode1]
volumes:
- ./conf:/etc/hadoop/conf:ro
- dn3-data:/hadoop/dfs/data
networks:
- hadoop-net
networks:
hadoop-net:
driver: bridge
volumes:
jn1-data: {}
jn2-data: {}
jn3-data: {}
nn1-data: {}
nn2-data: {}
dn1-data: {}
dn2-data: {}
dn3-data: {}
Лабораторное задание: симуляция failover под нагрузкой¶
# lab_failover_simulation.py
# Запускаем тяжёлый Spark-запрос и в середине выполнения убиваем Active NameNode.
# Цель: убедиться, что джоба завершается успешно несмотря на failover.
import subprocess
import threading
import time
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
def create_spark_ha() -> SparkSession:
"""
SparkSession с конфигурацией HA для Docker Compose стенда.
Конфиги прописаны явно (не через hdfs-site.xml) для простоты демонстрации.
"""
return (
SparkSession.builder
.master("local[4]")
.appName("failover-lab")
.config("spark.hadoop.fs.defaultFS", "hdfs://mycluster")
.config("spark.hadoop.dfs.nameservices", "mycluster")
.config("spark.hadoop.dfs.ha.namenodes.mycluster", "nn1,nn2")
.config("spark.hadoop.dfs.namenode.rpc-address.mycluster.nn1", "localhost:8020")
.config("spark.hadoop.dfs.namenode.rpc-address.mycluster.nn2", "localhost:8021")
.config("spark.hadoop.dfs.client.failover.proxy.provider.mycluster",
"org.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider")
# Таймауты адаптированы под Docker окружение
.config("spark.hadoop.dfs.client.failover.max.attempts", "10")
.config("spark.hadoop.dfs.client.failover.sleep.base.millis", "300")
.config("spark.hadoop.ipc.client.connect.timeout", "5000")
.getOrCreate()
)
def kill_active_namenode_delayed(delay_seconds: int = 30) -> threading.Thread:
"""
Убивает контейнер namenode1 через delay_seconds секунд.
Используем docker stop для имитации «мягкого» падения.
Для «жёсткого» используйте docker kill.
"""
def _kill():
print(f"\n[FAILOVER] Ждём {delay_seconds} секунд до симуляции аварии...")
time.sleep(delay_seconds)
print("[FAILOVER] Останавливаем namenode1 (docker stop)...")
result = subprocess.run(
["docker", "stop", "namenode1"],
capture_output=True, text=True
)
print(f"[FAILOVER] Результат: {result.returncode}")
t = threading.Thread(target=_kill, daemon=True)
t.start()
return t
def generate_test_dataset(spark: SparkSession, path: str, size_mb: int = 200) -> None:
"""
Генерирует тестовый датасет на HDFS перед запуском лабораторной работы.
size_mb: приблизительный размер датасета.
"""
rows_per_mb = 10_000
n_rows = size_mb * rows_per_mb
print(f"Генерируем датасет: {n_rows:,} строк → {path}")
df = spark.range(n_rows).select(
F.col("id").alias("event_id"),
(F.rand() * 1000).cast("long").alias("user_id"),
F.array(
F.lit("purchase"), F.lit("view"), F.lit("click"), F.lit("share")
).getItem((F.rand() * 4).cast("int")).alias("event_type"),
(F.rand() * 500).alias("amount"),
F.date_add(F.lit("2024-01-01"), (F.rand() * 365).cast("int")).alias("event_date"),
)
df.repartition(50).write.mode("overwrite").partitionBy("event_date").parquet(path)
print(f"Датасет создан: {n_rows:,} строк")
def run_heavy_aggregation(spark: SparkSession, input_path: str) -> dict:
"""
Запускает многошаговую агрегацию - достаточно долгую, чтобы
в процессе успело произойти переключение NameNode.
"""
print(f"\n[SPARK] Начинаем агрегацию из {input_path}")
start_time = time.time()
df = spark.read.parquet(input_path)
# Первая агрегация: по дате и типу события
daily_stats = df.groupBy("event_date", "event_type").agg(
F.count("*").alias("event_count"),
F.countDistinct("user_id").alias("unique_users"),
F.sum("amount").alias("total_amount"),
)
# Вторая агрегация поверх первой - добавляет Stage'ы и растягивает время
monthly_stats = daily_stats \
.withColumn("month", F.date_format("event_date", "yyyy-MM")) \
.groupBy("month", "event_type") \
.agg(
F.sum("event_count").alias("monthly_events"),
F.sum("unique_users").alias("monthly_users"),
F.sum("total_amount").alias("monthly_revenue"),
)
result_count = monthly_stats.count()
elapsed = time.time() - start_time
print(f"\n[SPARK] Агрегация завершена!")
print(f"[SPARK] Строк в результате: {result_count}")
print(f"[SPARK] Время выполнения: {elapsed:.1f}с")
return {"result_count": result_count, "elapsed_seconds": elapsed, "success": True}
def verify_failover() -> dict:
"""Проверяет состояние NameNode после теста."""
result = subprocess.run(
["docker", "exec", "namenode2", "hdfs", "haadmin", "-getServiceState", "nn2"],
capture_output=True, text=True
)
nn2_state = result.stdout.strip()
return {
"nn2_state": nn2_state,
"failover_detected": "active" in nn2_state.lower(),
}
if __name__ == "__main__":
spark = create_spark_ha()
INPUT_PATH = "hdfs://mycluster/lab/events/"
# Шаг 1: создаём тестовые данные (до аварии)
generate_test_dataset(spark, INPUT_PATH, size_mb=200)
# Шаг 2: планируем убийство Active NN через 30 секунд
kill_thread = kill_active_namenode_delayed(delay_seconds=30)
# Шаг 3: запускаем тяжёлую джобу (~60-90 секунд)
result = run_heavy_aggregation(spark, INPUT_PATH)
kill_thread.join(timeout=5)
# Шаг 4: верифицируем, что failover произошёл
status = verify_failover()
print("\n" + "=" * 60)
print("РЕЗУЛЬТАТЫ ЛАБОРАТОРНОЙ РАБОТЫ")
print("=" * 60)
print(f"Spark Job завершилась успешно: {result['success']}")
print(f"Время выполнения: {result['elapsed_seconds']:.1f}с")
print(f"NameNode 2 стал Active: {status['failover_detected']}")
print(f"Текущее состояние nn2: {status['nn2_state']}")
spark.stop()
Диагностика проблем HA: дерево решений¶
# Диагностические команды для типичных проблем
# 1. Проверить состояние всех NameNode
hdfs haadmin -getAllServiceState
# 2. Детальная информация о DataNode и свободном месте
hdfs dfsadmin -report
# 3. Состояние ZooKeeper lock
zkCli.sh -server zk1:2181,zk2:2181,zk3:2181 \
get /hadoop-ha/mycluster/ActiveStandbyElectorLock
# 4. Проверить JournalNode - активные сегменты EditLog
hdfs journalnode -listEditLogs
# 5. Логи ZKFC для анализа причины failover
tail -100f /var/log/hadoop/hdfs/hadoop-hdfs-zkfc-namenode1.log
# 6. Принудительный переход в Active (ОПАСНО)
# Только если уверены что другой NN мёртв
hdfs haadmin -transitionToActive --forceactive nn1
# 7. Проверить синхронизацию Standby
hdfs haadmin -checkHealth nn2
Итоги: что важно помнить о NameNode HA¶
NameNode HA - это не просто «резервный сервер». Это комплексная архитектура из нескольких взаимодействующих компонентов, каждый из которых решает конкретную проблему.
QJM / JournalNode решают задачу синхронизации метаданных между Active и Standby без риска потери данных. Механизм Epoch ID гарантирует, что только один NameNode может писать в журнал в каждый момент времени.
ZooKeeper + ZKFC решают задачу автоматического обнаружения отказов и координации выбора нового Active без участия администратора. Эфемерные znode - элегантное решение проблемы «кто владеет замком».
Fencing решает задачу предотвращения Split-Brain - самой опасной ситуации в HA-системах. Правило «сначала убить старый Active, потом активировать новый» кажется очевидным, но требует надёжного механизма принудительного отключения.
ConfiguredFailoverProxyProvider решает задачу прозрачности для клиентов: Spark-код пишет hdfs://mycluster/path и не задумывается о том, nn1 или nn2 сейчас Active. При failover клиент автоматически переключается.
Для Spark Data Engineer практические выводы:
- Используйте логическое имя nameservice (
hdfs://mycluster) вместо конкретного хоста во всех пайплайнах - Настройте правильные таймауты failover под характер джобов (стриминг vs batch)
- Плановое обслуживание NameNode делайте через
hdfs haadmin -failover, а не через прямойsystemctl stop - Фазы Spark Job с операциями над namespace (listStatus, create, rename) уязвимы к failover - проектируйте идемпотентные пайплайны
- Аномально высокий Scheduler Delay в Spark UI - это сигнатура failover-события: ищите его первым при анализе медленных Stage'ов