Short-Circuit Reads: чтение DataNode без сетевого стека через Unix socket
Глубокий разбор механизма Short-Circuit Local Reads в HDFS: физика оверхеда стандартного TCP-чтения, архитектура Unix Domain Socket и File Descriptor Passing, роль libhadoop.so и shared memory, пошаговая конфигурация hdfs-site.xml, тюнинг PySpark, мониторинг через JMX и траблшутинг типичных ошибок.
1. Физическая проблема стандартного NODE_LOCAL чтения в HDFS¶
В предыдущем уроке мы разобрали, что DATA Locality в HDFS - это механизм, позволяющий Spark-задачам читать данные с локального диска того же сервера, где запущен Executor. NODE_LOCAL кажется идеальным: данные рядом, сеть не нужна. Но на практике даже NODE_LOCAL чтение в стандартной конфигурации проходит через сетевой стек ядра Linux. Это и есть проблема, которую решает Short-Circuit.
Путь данных в стандартном NODE_LOCAL чтении¶
Когда Spark Executor запрашивает блок HDFS, который физически лежит на том же сервере, происходит следующее:
Схема раскрывает ключевую проблему: несмотря на то что Spark Executor и DataNode работают на одном физическом сервере, данные совершают длинный путь. Самое важное - это двойное копирование данных в памяти:
- Первое копирование: disk (или Page Cache) → DataNode Java Heap. DataNode читает блок в свой буфер (
byte[]в heap). - Второе копирование: DataNode Heap → TCP Socket Buffer (ядро) → Loopback → Executor Heap. Данные проходят через сетевой стек ядра, даже используя loopback интерфейс.
Анатомия оверхеда: что именно тратит CPU¶
Каждый раз, когда данные копируются через TCP loopback, Linux ядро выполняет несколько дорогостоящих операций:
Context Switches (переключения контекста). Процесс DataNode вызывает write() в TCP сокет - ядро переключается в kernel space. Процесс Spark Executor вызывает read() - снова kernel space. Для одного блока 128 MB при размере буфера 64 KB это означает 2048 переключений контекста только для одного блока (128 MB / 64 KB = 2048 операций write/read пар).
При чтении 100 блоков параллельно (нормальная нагрузка Spark Stage) - это 200 000 переключений контекста за несколько секунд. Это значительная нагрузка на CPU планировщик ядра.
Memory Copies. Linux TCP loopback не является нулевым копированием (zero-copy). Даже при localhost соединении 127.0.0.1 → 127.0.0.1 данные копируются: DataNode Heap → Socket Send Buffer → Socket Receive Buffer → Executor Heap. Итого: 3 копии данных в памяти для одного блока 128 MB = 384 MB выделено и скопировано только на передачу одного блока.
TCP Handshake и Protocol Overhead. Каждое открытие нового TCP соединения на DataNode требует трёхстороннего рукопожатия (SYN, SYN-ACK, ACK) и HDFS протокольного handshake - это десятки миллисекунд задержки перед началом фактической передачи данных.
Числа: насколько это серьёзно¶
Чтобы оценить масштаб проблемы, рассмотрим реальный сценарий: Spark Stage читает 500 GB данных в формате Parquet с 50 Executor'ов на кластере из 10 серверов, каждый из которых имеет DataNode.
Стандартное TCP чтение (NODE_LOCAL):
- 500 GB данных × 3 memory copies = ~1.5 TB переданных байт
- Context switches: (500 GB / 128 MB) × 2048 × 2 = ~8 млн переключений
- Итоговое время: ~80-120 секунд (ограничение CPU на context switches)
Short-Circuit чтение:
- 500 GB данных × 1 memory copy = ~500 GB переданных байт
- Context switches: (500 GB / 128 MB) × 4 = ~16000 переключений (в 500 раз меньше!)
- Итоговое время: ~30-50 секунд
Разница в 2–3 раза по общему времени выполнения только за счёт устранения сетевого стека на локальных чтениях. Именно за это Short-Circuit Reads называют «одним из главных скрытых резервов производительности HDFS-кластеров».
2. Концепция и механика Short-Circuit Local Reads¶
Short-Circuit (буквально «короткое замыкание») - это оптимизация, при которой HDFS Client Spark'а замыкает путь данных напрямую на локальную файловую систему, обходя DataNode как посредника передачи данных.
Термин «short-circuit» пришёл из электроники: когда электрический ток находит путь с минимальным сопротивлением (короткое замыкание), он идёт по нему вместо длинного основного пути. Аналогично: данные находят «короткий путь» disk → Executor heap, минуя DataNode heap и TCP stack.
Что именно делает Short-Circuit¶
При Short-Circuit чтении HDFS Client Spark'а запрашивает у DataNode не данные, а файловый дескриптор (File Descriptor) для файла блока на локальной файловой системе. Получив дескриптор, Executor сам открывает файл и читает данные напрямую:
Сравнение с предыдущей схемой показывает принципиальное упрощение пути данных: теперь от шага 7 (read()) данные попадают сразу в Executor Heap (шаг 11), минуя DataNode heap и TCP буферы. Одна копия вместо трёх.
Unix Domain Socket: почему это быстрее TCP loopback¶
Unix Domain Socket (UDS) - это механизм межпроцессного взаимодействия (IPC), встроенный в ядро Linux. В отличие от TCP sockets:
- Нет сетевого протокола: UDS не проходит через TCP/IP стек, нет заголовков пакетов, нет congestion control, нет TCP handshake
- Адресация по файловому пути: вместо
127.0.0.1:50010используется/var/lib/hadoop-hdfs/dn_socket- путь на файловой системе - Прямая передача в ядре: данные передаются через kernel buffer напрямую между двумя процессами без копирования в сетевой стек
- File Descriptor Passing: UDS поддерживает специальный механизм передачи файловых дескрипторов между процессами (SCM_RIGHTS) - именно это делает Short-Circuit возможным
File Descriptor Passing - это уникальная возможность Unix Domain Sockets, которой нет у TCP. Когда DataNode открывает файл блока (open() возвращает fd=7), DataNode может передать этот дескриптор другому процессу (Spark Executor) через UDS. Executor получает свою копию дескриптора (fd=9 - другой номер, но указывает на тот же файл в ядре) и может читать файл напрямую.
Почему это безопасно? DataNode контролирует доступ: он передаёт файловый дескриптор только после проверки Block Access Token - криптографического токена, который Spark получил от NameNode при запросе блока. Без валидного токена DataNode откажет в передаче дескриптора.
3. Разделяемая память: libhadoop.so и Shared Memory Segments¶
Short-Circuit работает не только через Unix Domain Socket. Есть ещё более продвинутый слой оптимизации - разделяемая память (Shared Memory) для хранения метаданных о состоянии блоков.
Зачем нужна разделяемая память при Short-Circuit¶
При чтении блока через Short-Circuit Spark Executor должен верифицировать контрольную сумму (checksum) прочитанных данных. HDFS хранит checksums в отдельном файле .meta рядом с каждым блоком:
/data/hdfs/current/BP-.../blk_001 ← сами данные блока
/data/hdfs/current/BP-.../blk_001.meta ← checksums для каждого чанка блока
Стандартный подход: Executor читает .meta файл и сверяет checksums. Но есть проблема: при интенсивном чтении Executor должен постоянно запрашивать у DataNode актуальное состояние блока (не устарел ли он, не помечен ли как corrupted). Каждый такой запрос - roundtrip через UDS.
Shared Memory Segment решает это: DataNode и HDFS Client Spark'а разделяют одну область памяти ядра, где DataNode пишет статус блоков, а Spark читает без каких-либо roundtrip запросов:
Когда DataNode помечает блок как STALE (например, идёт репликация или блок повреждён), Spark видит это мгновенно через shared memory и автоматически переключается на чтение через TCP от другого DataNode. Нет опроса, нет задержки обнаружения.
Нативная библиотека libhadoop.so: почему она критически важна¶
Работа с Unix Domain Sockets и Shared Memory требует системных вызовов, которые Java не может выполнить напрямую из-за ограничений JVM sandbox. Для этого HDFS использует нативную C++ библиотеку libhadoop.so.
libhadoop.so предоставляет через JNI (Java Native Interface) следующие возможности:
sendFd()/receiveFd()- передача файловых дескрипторов через UDS (SCM_RIGHTS). Java не поддерживает это нативно.mmap()- маппинг shared memory сегментов. Чистая Java не имеет прямого доступа к POSIX mmap.fadvise()- подсказки ядру о паттерне чтения файла (POSIX_FADV_SEQUENTIAL для Parquet sequential scan), что позволяет ядру оптимизировать prefetch.munlock()/mlock()- управление блокировкой страниц памяти для предотвращения свопирования горячих данных.
Что происходит без libhadoop.so? Если нативная библиотека не найдена или не загружена, HDFS Client автоматически отваливается на стандартное TCP чтение. Никакой ошибки не выдаётся - просто тихий fallback. Именно поэтому многие кластеры работают годами без Short-Circuit, думая что оно включено.
Признак отсутствия libhadoop в логах Spark Driver:
WARN NativeCodeLoader: Unable to load native-hadoop library for your platform...
using builtin-java classes where applicable
Это предупреждение означает: Short-Circuit не работает, всё читается через TCP, даже NODE_LOCAL задачи.
Где найти и как проверить libhadoop.so¶
# Проверить наличие нативных библиотек Hadoop
find /usr/lib/hadoop /opt/hadoop -name "libhadoop.so*" 2>/dev/null
# Ожидаемый результат: /usr/lib/hadoop/lib/native/libhadoop.so.1.0.0
# Проверить, что библиотека совместима с текущей ОС
file /usr/lib/hadoop/lib/native/libhadoop.so.1.0.0
# ELF 64-bit LSB shared object, x86-64, dynamically linked
# Проверить символы (убедиться что FD passing функции присутствуют)
nm -D /usr/lib/hadoop/lib/native/libhadoop.so | grep -i "sendfd\|receivefd\|mmap"
# Должны быть видны: Java_org_apache_hadoop_io_nativeio_NativeIO_00024POSIX_sendFileDescriptor
# Установить переменную окружения для Spark
export HADOOP_OPTS="-Djava.library.path=/usr/lib/hadoop/lib/native"
# Или в spark-env.sh:
echo 'export LD_LIBRARY_PATH=/usr/lib/hadoop/lib/native:$LD_LIBRARY_PATH' >> $SPARK_HOME/conf/spark-env.sh
4. Пошаговая конфигурация инфраструктуры: hdfs-site.xml¶
Конфигурация Short-Circuit Reads требует согласованных изменений на двух уровнях: серверном (DataNode + операционная система) и клиентском (Spark). Начнём с серверной стороны.
Подготовка операционной системы: создание директории сокета¶
# На каждом узле с DataNode:
# Создаём директорию для Unix Domain Socket
# Важно: не /tmp! Там могут создавать файлы другие процессы.
# Используем /var/lib/hadoop-hdfs - отдельная директория с жёсткими правами.
sudo mkdir -p /var/lib/hadoop-hdfs
# Устанавливаем владельца - пользователь под которым работает DataNode
# Обычно это hdfs:hadoop
sudo chown hdfs:hadoop /var/lib/hadoop-hdfs
# Права доступа КРИТИЧЕСКИ ВАЖНЫ для безопасности.
# 700: только hdfs (владелец) может читать/записывать/исполнять.
# Если поставить 777 - Spark откажется использовать Short-Circuit
# с ошибкой "The domain socket path is insecure"!
sudo chmod 700 /var/lib/hadoop-hdfs
# Проверяем результат
ls -la /var/lib/ | grep hadoop-hdfs
# drwx------ 2 hdfs hadoop 4096 Jan 15 10:00 hadoop-hdfs
Полная конфигурация hdfs-site.xml¶
<!-- hdfs-site.xml - применяется на DataNode серверах и Spark клиентах -->
<!-- ══════════════════════════════════════════════════════════
CORE Short-Circuit конфигурация
══════════════════════════════════════════════════════════ -->
<!--
dfs.client.read.shortcircuit: главный флаг включения Short-Circuit.
true: HDFS Client будет запрашивать File Descriptor у DataNode
вместо данных. Если DataNode на другом сервере или SC недоступен —
автоматический fallback на TCP без ошибок.
false (по умолчанию): всегда TCP, Short-Circuit не используется.
-->
<property>
<name>dfs.client.read.shortcircuit</name>
<value>true</value>
</property>
<!--
dfs.domain.socket.path: путь к Unix Domain Socket файлу.
DataNode создаёт этот файл при старте и слушает на нём.
Spark Client подключается к нему для получения File Descriptor.
ВАЖНО: директория должна принадлежать пользователю DataNode
с правами 700 (только владелец). Иначе Spark откажется
использовать SC из соображений безопасности.
Рекомендуемый путь: /var/lib/hadoop-hdfs/dn_socket
Не используйте /tmp - там нет контроля прав!
-->
<property>
<name>dfs.domain.socket.path</name>
<value>/var/lib/hadoop-hdfs/dn_socket</value>
</property>
<!-- ══════════════════════════════════════════════════════════
Checksum и целостность данных
══════════════════════════════════════════════════════════ -->
<!--
dfs.client.read.shortcircuit.skip.checksum: пропускать проверку
контрольных сумм при Short-Circuit чтении.
false (рекомендуется для production): проверять checksums.
Это защищает от silent data corruption (bit rot).
При Short-Circuit checksum проверяется локально из .meta файла
без обращения к DataNode - overhead минимален.
true: пропустить checksum для максимальной скорости.
Допустимо ТОЛЬКО для временных/тестовых данных, где потеря
части данных некритична. Никогда не использовать на production
financial или medical данных.
-->
<property>
<name>dfs.client.read.shortcircuit.skip.checksum</name>
<value>false</value>
</property>
<!-- ══════════════════════════════════════════════════════════
Shared Memory и Block State Management
══════════════════════════════════════════════════════════ -->
<!--
dfs.client.read.shortcircuit.streams.cache.size: размер кеша
открытых файловых дескрипторов на клиенте.
При Short-Circuit для каждого блока открывается fd.
Кеш позволяет переиспользовать fd для повторных чтений
того же блока (например при partition pruning нескольких запросов).
Дефолт: 256. Увеличьте до 1000-2000 для больших Stage'ов
с высоким параллелизмом (200+ одновременных задач).
Каждый открытый fd потребляет ~1 KB памяти.
-->
<property>
<name>dfs.client.read.shortcircuit.streams.cache.size</name>
<value>1000</value>
</property>
<!--
dfs.client.read.shortcircuit.streams.cache.expiry.ms: время
жизни кешированного fd без использования.
Дефолт: 300000 мс (5 минут). Fd автоматически закрывается
если не использовался указанное время.
Снижайте для кластеров с большим числом файлов и высоким
file descriptor pressure (ulimit -n).
Увеличивайте для долгих джобов с повторными сканированиями.
-->
<property>
<name>dfs.client.read.shortcircuit.streams.cache.expiry.ms</name>
<value>600000</value>
</property>
<!--
dfs.short.circuit.shared.memory.watcher.interrupt.check.ms:
как часто ShortCircuitShm менеджер проверяет состояние
shared memory сегментов.
Дефолт: 60000 мс (1 минута). При агрессивном перераспределении
блоков на кластере можно снизить до 10000-30000 мс для
более быстрого обнаружения устаревших блоков.
-->
<property>
<name>dfs.short.circuit.shared.memory.watcher.interrupt.check.ms</name>
<value>30000</value>
</property>
<!-- ══════════════════════════════════════════════════════════
Разрешения доступа
══════════════════════════════════════════════════════════ -->
<!--
dfs.block.local-path-access.user: список пользователей,
которым DataNode разрешает Short-Circuit доступ к блокам.
Пользователь, от имени которого запускается Spark (yarn, spark,
mapred) должен быть в этом списке. Иначе DataNode отклонит
запрос на File Descriptor с ошибкой AccessControlException.
Пример: "spark,yarn,mapred,hdfs"
В Kerberos окружениях используйте полные principal имена.
-->
<property>
<name>dfs.block.local-path-access.user</name>
<value>spark,yarn,mapred</value>
</property>
<!--
dfs.datanode.data.dir.perm: права доступа к директориям
с данными DataNode. Минимум 750 (r-xr-x---):
DataNode (владелец) - полный доступ
hadoop group - чтение и исполнение (для Short-Circuit)
Другие - нет доступа
755 даёт более широкий доступ, что упрощает настройку
но снижает безопасность для мультитенантных кластеров.
-->
<property>
<name>dfs.datanode.data.dir.perm</name>
<value>750</value>
</property>
Верификация конфигурации после применения¶
# Перезапускаем DataNode для применения конфигурации
sudo systemctl restart hadoop-hdfs-datanode
# Проверяем что DataNode создал Unix socket файл
ls -la /var/lib/hadoop-hdfs/dn_socket
# srwxr-xr-x 1 hdfs hadoop 0 Jan 15 10:05 /var/lib/hadoop-hdfs/dn_socket
# Тип файла 's' = socket (не обычный файл!)
# Убеждаемся что DataNode слушает на сокете
sudo -u hdfs hadoop dfsadmin -report
# В выводе должно быть: "Configured Capacity: ... DFS Used: ..."
# Проверяем через lsof что DataNode открыл Unix socket
sudo lsof -U | grep dn_socket
# hadoop.da 12345 hdfs xxx unix /var/lib/hadoop-hdfs/dn_socket type=STREAM
# Тест: попробуем подключиться к сокету как пользователь spark
sudo -u spark python3 -c "
import socket
s = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
try:
s.connect('/var/lib/hadoop-hdfs/dn_socket')
print('SUCCESS: пользователь spark может подключиться к UDS')
s.close()
except PermissionError as e:
print(f'ОШИБКА прав доступа: {e}')
"
5. Проброс настроек и тюнинг на стороне PySpark¶
Серверная конфигурация в hdfs-site.xml - это только половина настройки. Spark должен знать о Short-Circuit и правильно передать конфигурацию всем компонентам: Driver и Executor'ам.
Автоматическая загрузка конфигурации из hdfs-site.xml¶
Если Spark запускается на кластере с правильно настроенным HADOOP_CONF_DIR, он автоматически загружает core-site.xml и hdfs-site.xml:
# hadoop-env.sh или spark-env.sh
export HADOOP_CONF_DIR=/etc/hadoop/conf
# Spark читает все XML файлы в этой директории при старте
# и применяет их к Hadoop Configuration объекту внутри SparkContext
Однако полагаться только на автозагрузку рискованно: разные машины кластера могут иметь разные версии конфига. Рекомендуется явно дублировать критичные параметры в SparkSession.
Полная конфигурация SparkSession для Short-Circuit¶
# spark_shortcircuit_config.py
import os
from pyspark.sql import SparkSession
def create_spark_with_shortcircuit(
app_name: str,
nameservice: str = "mycluster",
socket_path: str = "/var/lib/hadoop-hdfs/dn_socket",
enable_shortcircuit: bool = True,
fd_cache_size: int = 1000,
) -> SparkSession:
"""
Создаёт SparkSession с явной конфигурацией Short-Circuit Reads.
Параметры:
app_name: имя Spark-приложения
nameservice: логическое имя HDFS HA nameservice
socket_path: путь к Unix Domain Socket DataNode
enable_shortcircuit: включить или отключить SC (для A/B тестирования)
fd_cache_size: размер кеша файловых дескрипторов на Executor
"""
builder = SparkSession.builder \
.appName(app_name) \
# ── HDFS HA конфигурация ─────────────────────────────────────
.config("spark.hadoop.fs.defaultFS", f"hdfs://{nameservice}") \
# ── Short-Circuit Reads ───────────────────────────────────────
# Все параметры с префиксом spark.hadoop. передаются в
# Hadoop Configuration объект на Driver и каждом Executor.
# Это эквивалентно записи в hdfs-site.xml, но на уровне сессии.
.config("spark.hadoop.dfs.client.read.shortcircuit",
"true" if enable_shortcircuit else "false") \
# Путь к UDS файлу DataNode.
# Должен совпадать с dfs.domain.socket.path на DataNode!
.config("spark.hadoop.dfs.domain.socket.path", socket_path) \
# Выполнять checksum проверку (true = безопасно, рекомендуется)
.config("spark.hadoop.dfs.client.read.shortcircuit.skip.checksum",
"false") \
# Кеш файловых дескрипторов.
# Увеличьте если кластер большой и много параллельных задач.
.config("spark.hadoop.dfs.client.read.shortcircuit.streams.cache.size",
str(fd_cache_size)) \
# Время жизни кешированного fd (10 минут)
.config("spark.hadoop.dfs.client.read.shortcircuit.streams.cache.expiry.ms",
"600000") \
# ── Нативные библиотеки ───────────────────────────────────────
# Путь к libhadoop.so для JNI.
# Без этого Short-Circuit не работает - будет тихий fallback на TCP!
# Передаём в JVM аргументы Driver и Executor'ов.
.config("spark.driver.extraJavaOptions",
"-Djava.library.path=/usr/lib/hadoop/lib/native") \
.config("spark.executor.extraJavaOptions",
"-Djava.library.path=/usr/lib/hadoop/lib/native") \
# Также нужно чтобы ОС нашла libhadoop.so через LD_LIBRARY_PATH
# Это особенно важно для YARN контейнеров
.config("spark.executorEnv.LD_LIBRARY_PATH",
"/usr/lib/hadoop/lib/native") \
.config("spark.driverEnv.LD_LIBRARY_PATH",
"/usr/lib/hadoop/lib/native") \
# ── Data Locality (связано с Short-Circuit) ───────────────────
# NODE_LOCAL - наиболее критичный уровень для SC.
# Увеличиваем wait, чтобы задачи чаще получали NODE_LOCAL.
.config("spark.locality.wait.node", "10s") \
.config("spark.locality.wait.process", "0") \
# ── Производительность I/O ────────────────────────────────────
# Размер буфера чтения. 128 KB оптимален для Parquet Row Groups.
# При Short-Circuit данные читаются большими буферами с диска.
.config("spark.hadoop.io.file.buffer.size", "131072") \
# Prefetching: сколько блоков Spark заранее запрашивает
# у NameNode для prefetch. При SC это менее важно (нет сетевого
# round-trip), но улучшает pipeline при последовательном сканировании.
.config("spark.hadoop.dfs.client.read.shortcircuit.buffer.size",
"1048576") # 1 MB буфер для SC чтения
spark = builder.getOrCreate()
# Устанавливаем уровень логирования для диагностики SC
# DEBUG покажет каждый SC-запрос, INFO - только проблемы
spark.sparkContext.setLogLevel("WARN")
# Включаем подробные логи SC если нужна диагностика
# spark.sparkContext._jvm.org.apache.log4j.Logger \
# .getLogger("org.apache.hadoop.hdfs.BlockReaderLocal") \
# .setLevel(spark.sparkContext._jvm.org.apache.log4j.Level.DEBUG)
return spark
# ── Примеры использования ────────────────────────────────────────────
# Стандартная production конфигурация
spark_prod = create_spark_with_shortcircuit(
app_name="etl-bronze-ingestion",
nameservice="mycluster",
socket_path="/var/lib/hadoop-hdfs/dn_socket",
enable_shortcircuit=True,
fd_cache_size=2000, # Для больших Stage'ов
)
# Для A/B теста: SC выключен
spark_no_sc = create_spark_with_shortcircuit(
app_name="etl-no-shortcircuit-baseline",
enable_shortcircuit=False,
)
Передача конфигурации через spark-defaults.conf¶
Для production кластеров удобнее задавать Short-Circuit конфигурацию глобально, чтобы все Spark-приложения получали её автоматически:
# /etc/spark/conf/spark-defaults.conf
# Применяется ко всем Spark-приложениям на кластере
# Short-Circuit Reads
spark.hadoop.dfs.client.read.shortcircuit true
spark.hadoop.dfs.domain.socket.path /var/lib/hadoop-hdfs/dn_socket
spark.hadoop.dfs.client.read.shortcircuit.skip.checksum false
spark.hadoop.dfs.client.read.shortcircuit.streams.cache.size 1000
# Нативные библиотеки
spark.driver.extraJavaOptions -Djava.library.path=/usr/lib/hadoop/lib/native
spark.executor.extraJavaOptions -Djava.library.path=/usr/lib/hadoop/lib/native
spark.executorEnv.LD_LIBRARY_PATH /usr/lib/hadoop/lib/native
# Data Locality (комплементарная настройка)
spark.locality.wait.node 10s
spark.locality.wait.process 0
Права доступа JVM-процесса: пользователи и группы Linux¶
Критически важный и часто упускаемый аспект - POSIX права пользователя, под которым запускается Spark Executor, должны позволять ему:
- Подключиться к Unix Domain Socket DataNode
- Прочитать файлы блоков через полученный File Descriptor
# Сценарий настройки прав для пользователя spark
# 1. Убеждаемся что пользователь spark существует
id spark
# uid=999(spark) gid=999(spark) groups=999(spark)
# 2. Добавляем spark в группу hadoop для доступа к data директориям DataNode
sudo usermod -aG hadoop spark
# 3. Проверяем права на директорию сокета
ls -la /var/lib/ | grep hadoop-hdfs
# drwx--x--- 2 hdfs hadoop 4096 Jan 15 10:00 hadoop-hdfs
# ^^^^^^^
# hadoop group может исполнять (traverse) директорию - это нужно для connect()
# 4. Проверяем права на data директории DataNode
# Файлы блоков должны быть читаемы для группы hadoop
ls -la /data/hdfs/current/BP-*/
# -rw-r----- 1 hdfs hadoop 134217728 Jan 15 08:00 blk_001
# -rw-r----- 1 hdfs hadoop 1052678 Jan 15 08:00 blk_001.meta
# Права 640 или 644 - оба подходят
# 5. Если дата директории имеют права 700 (только hdfs)
# нужно изменить на 750 или использовать dfs.datanode.data.dir.perm=750
sudo find /data/hdfs -name "blk_*" | head -5 | xargs ls -la
# Если только hdfs может читать - spark получит Permission Denied при SC!
# 6. Применяем правильные права (если нужно изменить)
# Осторожно: это изменение влияет на все файлы DataNode!
# Лучше настроить через dfs.datanode.data.dir.perm в hdfs-site.xml
sudo find /data/hdfs -type f -name "blk_*" -exec chmod 640 {} \;
sudo find /data/hdfs -type d -exec chmod 750 {} \;
6. Мониторинг, аудит и верификация Short-Circuit¶
Включить конфигурацию мало - нужно убедиться, что Short-Circuit реально используется. Это не очевидно: при любой проблеме Spark тихо падает на TCP, не выдавая ошибки.
JMX-метрики DataNode: главный источник правды¶
DataNode публикует метрики Short-Circuit через JMX. Это единственный достоверный способ убедиться, что SC работает:
# Получаем метрики DataNode через HTTP JMX API
# Порт 9864 для HDFS 3.x (50075 для HDFS 2.x)
curl -s "http://datanode1.example.com:9864/jmx?qry=Hadoop:service=DataNode,name=DataNodeActivity" \
| python3 -c "
import sys, json
data = json.load(sys.stdin)
for bean in data['beans']:
if 'DataNodeActivity' not in bean.get('name', ''):
continue
# Основные метрики чтения
total_blocks_read = bean.get('BlocksRead', 0)
sc_reads = bean.get('ShortCircuitLocalReadsRead', 0)
sc_bytes = bean.get('ShortCircuitLocalBytesRead', 0)
tcp_reads = total_blocks_read - sc_reads
sc_pct = (sc_reads / total_blocks_read * 100) if total_blocks_read > 0 else 0
print(f'=== DataNode Short-Circuit Metrics ===')
print(f'Всего блоков прочитано: {total_blocks_read:,}')
print(f'Через Short-Circuit: {sc_reads:,} ({sc_pct:.1f}%)')
print(f'Через TCP fallback: {tcp_reads:,} ({100-sc_pct:.1f}%)')
print(f'Байт через SC: {sc_bytes / 1024**3:.2f} GB')
print()
# FD Cache статистика
fd_cache_size = bean.get('ShortCircuitLocalReadFdsAvg', 0)
fd_cache_misses = bean.get('ShortCircuitLocalReadFdsMissCount', 0)
print(f'=== FD Cache Metrics ===')
print(f'Avg открытых fd: {fd_cache_size}')
print(f'Cache misses: {fd_cache_misses:,}')
# Если cache miss высокий - увеличить streams.cache.size!
if fd_cache_misses > sc_reads * 0.1:
print(f'ВНИМАНИЕ: Cache miss rate > 10%! Увеличьте streams.cache.size')
"
Интерпретация результатов:
=== Здоровое состояние Short-Circuit ===
Всего блоков прочитано: 15234
Через Short-Circuit: 14891 (97.7%) ← отлично!
Через TCP fallback: 343 (2.3%) ← небольшой процент нормален (RACK_LOCAL задачи)
Байт через SC: 1.82 GB
=== Проблема: Short-Circuit не работает ===
Всего блоков прочитано: 15234
Через Short-Circuit: 0 (0.0%) ← Short-Circuit не работает!
Через TCP fallback: 15234 (100%) ← всё через TCP
Анализ метрик через Python и Prometheus¶
# sc_metrics_monitor.py
# Мониторинг Short-Circuit эффективности через JMX API
import requests
import json
from dataclasses import dataclass
from typing import Optional
@dataclass
class ShortCircuitMetrics:
hostname: str
total_blocks_read: int
sc_blocks_read: int
sc_bytes_read: int
tcp_blocks_read: int
@property
def sc_percentage(self) -> float:
if self.total_blocks_read == 0:
return 0.0
return self.sc_blocks_read / self.total_blocks_read * 100
@property
def is_healthy(self) -> bool:
"""SC считается здоровым если > 80% локальных блоков читается через SC."""
return self.sc_percentage > 80.0
def collect_sc_metrics(datanode_host: str, port: int = 9864) -> Optional[ShortCircuitMetrics]:
"""
Собирает Short-Circuit метрики с DataNode через JMX HTTP API.
Возвращает None если DataNode недоступен.
"""
url = f"http://{datanode_host}:{port}/jmx?qry=Hadoop:service=DataNode,name=DataNodeActivity"
try:
response = requests.get(url, timeout=5)
data = response.json()
for bean in data.get("beans", []):
if "DataNodeActivity" not in bean.get("name", ""):
continue
total = bean.get("BlocksRead", 0)
sc = bean.get("ShortCircuitLocalReadsRead", 0)
sc_bytes = bean.get("ShortCircuitLocalBytesRead", 0)
return ShortCircuitMetrics(
hostname=datanode_host,
total_blocks_read=total,
sc_blocks_read=sc,
sc_bytes_read=sc_bytes,
tcp_blocks_read=total - sc,
)
except (requests.RequestException, KeyError, json.JSONDecodeError) as e:
print(f"Не удалось получить метрики с {datanode_host}: {e}")
return None
def audit_cluster_sc_health(datanode_hosts: list[str]) -> dict:
"""
Аудит состояния Short-Circuit на всём кластере.
Возвращает сводку с проблемными узлами.
"""
results = []
problem_nodes = []
for host in datanode_hosts:
metrics = collect_sc_metrics(host)
if metrics is None:
problem_nodes.append({"host": host, "issue": "недоступен"})
continue
results.append(metrics)
if not metrics.is_healthy:
problem_nodes.append({
"host": host,
"issue": f"SC rate {metrics.sc_percentage:.1f}% < 80%",
"sc_blocks": metrics.sc_blocks_read,
"total_blocks": metrics.total_blocks_read,
})
total_sc = sum(m.sc_blocks_read for m in results)
total_reads = sum(m.total_blocks_read for m in results)
cluster_sc_pct = (total_sc / total_reads * 100) if total_reads > 0 else 0
print(f"\n{'='*50}")
print("АУДИТ SHORT-CIRCUIT НА КЛАСТЕРЕ")
print(f"{'='*50}")
print(f"Проверено DataNode: {len(results)}")
print(f"Кластерный SC rate: {cluster_sc_pct:.1f}%")
print(f"Всего блоков прочитано: {total_reads:,}")
print(f"Через SC: {total_sc:,}")
if problem_nodes:
print(f"\nПРОБЛЕМНЫЕ УЗЛЫ ({len(problem_nodes)}):")
for node in problem_nodes:
print(f" {node['host']}: {node['issue']}")
else:
print("\n✅ Все узлы работают корректно")
return {
"cluster_sc_percentage": cluster_sc_pct,
"problem_nodes": problem_nodes,
"healthy": len(problem_nodes) == 0,
}
# Пример использования
DATANODES = [
"datanode1.example.com",
"datanode2.example.com",
"datanode3.example.com",
]
audit_result = audit_cluster_sc_health(DATANODES)
Логи HDFS Client: что видно на стороне Spark¶
При включении DEBUG логирования для HDFS Client видны детали каждого SC-запроса:
# Включение DEBUG логов для HDFS Client в Spark (временно для диагностики)
spark-submit \
--conf "spark.driver.extraJavaOptions=-Dlog4j.logger.org.apache.hadoop.hdfs=DEBUG" \
--conf "spark.executor.extraJavaOptions=-Dlog4j.logger.org.apache.hadoop.hdfs=DEBUG" \
your_job.py
Характерные строки в логах при успешном Short-Circuit:
# Успешное подключение к Unix Domain Socket
DEBUG DomainSocketFactory: Successfully created Unix Domain Socket connection to
/var/lib/hadoop-hdfs/dn_socket
# Успешный Short-Circuit запрос к DataNode
DEBUG BlockReaderLocal: Creating new ShortCircuitReplicaInfo for
block BP-123456789-192.168.1.1-123/blk_001 on datanode1.example.com
# FD получен, читаем напрямую
DEBUG BlockReaderLocal: Successfully opened a short-circuit reader for
block BP-123456789-192.168.1.1-123/blk_001_1234.meta
# Чтение завершено
DEBUG BlockReaderLocal: read 134217728 bytes via short-circuit
Сигналы проблем с Short-Circuit:
# Ошибка прав доступа на socket (самая частая проблема!)
WARN DomainSocketFactory: error creating DomainSocket
/var/lib/hadoop-hdfs/dn_socket: java.io.IOException:
The UNIX Domain Socket at /var/lib/hadoop-hdfs/dn_socket is insecure.
It is owned by uid=0 (root), not uid=1000 (hdfs).
# После этого: тихий fallback на TCP
DEBUG BlockReaderRemote: Using remote blockReader to access
block BP-123456789.../blk_001 on datanode1:50010
# Нет libhadoop.so
WARN NativeCodeLoader: Unable to load native-hadoop library for your platform...
using builtin-java classes where applicable
# (Short-Circuit ПОЛНОСТЬЮ не работает без libhadoop!)
# FD Cache переполнен (нужно увеличить streams.cache.size)
DEBUG ShortCircuitCache: ShortCircuitReplica for block BP-.../blk_001 has expired.
streams.cache.size=256, current size=256
Влияние Short-Circuit на метрики Spark UI¶
При включённом Short-Circuit в Spark UI изменяются несколько метрик:
Task Deserialization Time снижается. Без SC каждый блок проходит через HDFS протокол de/serialization (ProtoBuf/ProtocolBuffer заголовки, RPC overhead). С SC этого нет - чистое чтение байт из файла.
GC Time снижается. Без SC DataNode создаёт промежуточные byte[] буферы в heap для каждого TCP пакета, Executor тоже создаёт буферы для принятых данных. Это значительная нагрузка на GC. С SC промежуточных объектов меньше - меньше GC давление.
Input Metrics → Bytes Read не меняется (данные читаются те же). Но Shuffle Read Time и Task Duration снижаются для Stage'ов с NODE_LOCAL задачами.
7. Практика: сравнительный бенчмарк и траблшутинг¶
Бенчмарк: Short-Circuit vs TCP на реальных данных¶
# benchmark_shortcircuit.py
# Сравниваем производительность чтения с Short-Circuit включённым и выключенным.
# Запускать на HDFS кластере с локальными DataNode!
import time
import subprocess
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
def create_spark(app_suffix: str, enable_sc: bool) -> SparkSession:
"""Создаёт SparkSession с включённым или выключенным Short-Circuit."""
sc_value = "true" if enable_sc else "false"
return (
SparkSession.builder
.master("yarn")
.appName(f"sc-benchmark-{'enabled' if enable_sc else 'disabled'}-{app_suffix}")
.config("spark.executor.instances", "10")
.config("spark.executor.cores", "4")
.config("spark.executor.memory", "8g")
# Short-Circuit конфигурация
.config("spark.hadoop.dfs.client.read.shortcircuit", sc_value)
.config("spark.hadoop.dfs.domain.socket.path",
"/var/lib/hadoop-hdfs/dn_socket")
.config("spark.hadoop.dfs.client.read.shortcircuit.skip.checksum", "false")
# Нативные библиотеки
.config("spark.executor.extraJavaOptions",
"-Djava.library.path=/usr/lib/hadoop/lib/native")
.config("spark.executorEnv.LD_LIBRARY_PATH",
"/usr/lib/hadoop/lib/native")
# Фиксируем locality.wait чтобы он не влиял на A/B тест
.config("spark.locality.wait.node", "10s")
# Отключаем AQE чтобы планы были идентичны
.config("spark.sql.adaptive.enabled", "false")
.getOrCreate()
)
def heavy_analytical_query(spark: SparkSession, data_path: str) -> tuple[int, float]:
"""
Тяжёлый аналитический запрос: полное сканирование крупного датасета.
Включает фильтрацию, агрегацию, join - типичный production ETL.
Возвращает: (число строк в результате, время выполнения в секундах)
"""
start = time.time()
# Читаем крупный Parquet датасет (ориентируемся на 50-100 GB)
events = spark.read.parquet(f"{data_path}/events/")
products = spark.read.parquet(f"{data_path}/products/")
# Типичный ETL запрос: фильтр + join + агрегация
result = (
events
.filter(
(F.col("event_type") == "purchase") &
(F.col("amount") > 10.0)
)
.join(
products.select("product_id", "category", "brand"),
"product_id",
"left"
)
.groupBy("category", "brand", F.date_format("event_time", "yyyy-MM-dd").alias("day"))
.agg(
F.count("*").alias("purchases"),
F.sum("amount").alias("revenue"),
F.countDistinct("user_id").alias("buyers"),
F.avg("amount").alias("avg_order_value"),
)
.orderBy(F.desc("revenue"))
)
count = result.count()
elapsed = time.time() - start
return count, elapsed
def get_context_switches() -> dict:
"""
Читает статистику context switches из /proc/stat.
Используется для измерения CPU overhead.
"""
with open("/proc/stat") as f:
for line in f:
if line.startswith("ctxt"):
return {"context_switches": int(line.split()[1])}
return {"context_switches": 0}
def run_full_benchmark(data_path: str) -> None:
"""Запускает полный A/B бенчмарк с Short-Circuit включённым и выключенным."""
results = {}
for sc_enabled in [False, True]:
label = "WITH Short-Circuit" if sc_enabled else "WITHOUT Short-Circuit (TCP)"
print(f"\n{'='*60}")
print(f"Запуск: {label}")
print(f"{'='*60}")
spark = create_spark(
app_suffix=str(int(time.time())),
enable_sc=sc_enabled
)
# Прогревочный запуск (warm up JVM, Parquet metadata cache)
print("Прогрев (warm-up run)...")
spark.read.parquet(f"{data_path}/events/").count()
# Измеряем context switches до теста
cs_before = get_context_switches()["context_switches"]
# Основной бенчмарк: 3 запуска для стабильности
times = []
for run in range(1, 4):
count, elapsed = heavy_analytical_query(spark, data_path)
times.append(elapsed)
print(f" Запуск {run}: {elapsed:.1f} сек, {count:,} строк")
# Измеряем context switches после теста
cs_after = get_context_switches()["context_switches"]
avg_time = sum(times) / len(times)
min_time = min(times)
results[label] = {
"avg_time_sec": avg_time,
"min_time_sec": min_time,
"context_switches_delta": cs_after - cs_before,
}
print(f"\nСредн. время: {avg_time:.1f} сек")
print(f"Мин. время: {min_time:.1f} сек")
print(f"Context switches за тест: {cs_after - cs_before:,}")
spark.stop()
time.sleep(10) # Пауза между тестами
# Сравнение результатов
print(f"\n{'='*60}")
print("ИТОГОВОЕ СРАВНЕНИЕ")
print(f"{'='*60}")
without_sc = results["WITHOUT Short-Circuit (TCP)"]
with_sc = results["WITH Short-Circuit"]
speedup = without_sc["avg_time_sec"] / with_sc["avg_time_sec"]
cs_reduction = (1 - with_sc["context_switches_delta"] /
without_sc["context_switches_delta"]) * 100
print(f"Время БЕЗ SC: {without_sc['avg_time_sec']:.1f} сек")
print(f"Время С SC: {with_sc['avg_time_sec']:.1f} сек")
print(f"Ускорение: {speedup:.2f}x")
print(f"Context switches снижены на {cs_reduction:.0f}%")
if __name__ == "__main__":
run_full_benchmark("hdfs://mycluster/benchmark/")
Траблшутинг: намеренно сломанная конфигурация и диагностика¶
Это самая практически ценная часть урока. Short-Circuit может не работать по десятку разных причин, и все они вызывают тихий fallback на TCP без явной ошибки в бизнес-логике джобы.
# troubleshooting_shortcircuit.py
# Диагностический скрипт для проверки Short-Circuit конфигурации.
# Запускать ОТ ИМЕНИ того пользователя, под которым работает Spark!
import os
import socket
import subprocess
import ctypes
from pathlib import Path
class ShortCircuitDiagnostics:
"""
Комплексная диагностика конфигурации Short-Circuit Reads.
Проверяет все необходимые условия и выдаёт конкретные рекомендации.
"""
def __init__(self, socket_path: str = "/var/lib/hadoop-hdfs/dn_socket",
hadoop_native_path: str = "/usr/lib/hadoop/lib/native"):
self.socket_path = socket_path
self.native_path = hadoop_native_path
self.issues = []
self.warnings = []
def check_socket_exists(self) -> bool:
"""Проверяет, существует ли файл Unix Domain Socket."""
path = Path(self.socket_path)
if not path.exists():
self.issues.append(
f"ОШИБКА: Unix Domain Socket не существует: {self.socket_path}\n"
f" Причина: DataNode не запущен или конфигурация dfs.domain.socket.path неверная\n"
f" Решение: запустить DataNode и проверить hdfs-site.xml"
)
return False
stat = path.stat()
# Проверяем что это именно socket (тип S_IFSOCK = 0o140000)
import stat as stat_module
if not stat_module.S_ISSOCK(stat.st_mode):
self.issues.append(
f"ОШИБКА: {self.socket_path} существует но не является Unix socket!\n"
f" Тип файла: {oct(stat.st_mode)}\n"
f" Решение: удалить файл и перезапустить DataNode"
)
return False
print(f"✅ Unix Domain Socket существует: {self.socket_path}")
return True
def check_socket_permissions(self) -> bool:
"""
Проверяет права доступа на директорию сокета.
Директория должна принадлежать пользователю DataNode (hdfs)
с правами максимум 755. Права 777 запрещены из соображений безопасности!
"""
socket_dir = Path(self.socket_path).parent
stat = socket_dir.stat()
import stat as stat_module
mode = stat.st_mode
perm_bits = stat_module.S_IMODE(mode)
# HDFS проверяет: если other (o) имеет write permission - это insecure!
if perm_bits & stat_module.S_IWOTH:
self.issues.append(
f"ОШИБКА БЕЗОПАСНОСТИ: директория {socket_dir} имеет write-доступ для 'others' "
f"(права: {oct(perm_bits)})\n"
f" Spark откажется использовать SC с ошибкой 'insecure domain socket path'\n"
f" Решение: sudo chmod o-w {socket_dir}\n"
f" Правильные права: 700 или 750 (не 777, не 775, не 755 с o+w)"
)
return False
# Проверяем что владелец директории - DataNode пользователь
import pwd
try:
owner = pwd.getpwuid(stat.st_uid).pw_name
if owner not in ("hdfs", "hadoop"):
self.warnings.append(
f"ПРЕДУПРЕЖДЕНИЕ: {socket_dir} принадлежит '{owner}', не 'hdfs'\n"
f" Может вызвать проблемы в некоторых конфигурациях"
)
except KeyError:
pass
print(f"✅ Права директории корректны: {oct(perm_bits)}")
return True
def check_socket_connectable(self) -> bool:
"""Пробует подключиться к Unix Domain Socket."""
s = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
s.settimeout(2.0)
try:
s.connect(self.socket_path)
s.close()
print(f"✅ Успешное подключение к Unix Domain Socket")
return True
except PermissionError:
current_user = os.getlogin()
self.issues.append(
f"ОШИБКА: Пользователь '{current_user}' не может подключиться к {self.socket_path}\n"
f" Решение 1: добавить пользователя в группу hadoop: usermod -aG hadoop {current_user}\n"
f" Решение 2: добавить пользователя в dfs.block.local-path-access.user в hdfs-site.xml"
)
return False
except ConnectionRefusedError:
self.issues.append(
f"ОШИБКА: DataNode не слушает на {self.socket_path}\n"
f" Причина: DataNode запущен но SC не настроен, или DataNode упал\n"
f" Проверьте: sudo systemctl status hadoop-hdfs-datanode"
)
return False
except socket.timeout:
self.issues.append(f"ОШИБКА: Таймаут подключения к {self.socket_path}")
return False
def check_libhadoop(self) -> bool:
"""Проверяет наличие и загружаемость нативной библиотеки Hadoop."""
libhadoop_files = list(Path(self.native_path).glob("libhadoop.so*"))
if not libhadoop_files:
self.issues.append(
f"ОШИБКА: libhadoop.so не найдена в {self.native_path}\n"
f" Short-Circuit НЕВОЗМОЖЕН без нативной библиотеки!\n"
f" Решение 1: установить hadoop-native пакет\n"
f" Решение 2: собрать Hadoop из исходников с -Pnative\n"
f" Решение 3: скачать pre-built бинарники под вашу ОС"
)
return False
# Пробуем загрузить библиотеку через ctypes
libhadoop_path = str(libhadoop_files[0])
try:
lib = ctypes.CDLL(libhadoop_path)
print(f"✅ libhadoop.so найдена и загружается: {libhadoop_path}")
return True
except OSError as e:
self.issues.append(
f"ОШИБКА: libhadoop.so найдена ({libhadoop_path}) но не загружается: {e}\n"
f" Причина: несовместимость с текущей ОС/архитектурой, или отсутствуют зависимости\n"
f" Диагностика: ldd {libhadoop_path}"
)
return False
def check_ld_library_path(self) -> bool:
"""Проверяет что LD_LIBRARY_PATH настроен для нахождения нативных библиотек."""
ld_path = os.environ.get("LD_LIBRARY_PATH", "")
java_lib_path = os.environ.get("JAVA_LIBRARY_PATH", "")
if self.native_path not in ld_path and self.native_path not in java_lib_path:
self.warnings.append(
f"ПРЕДУПРЕЖДЕНИЕ: {self.native_path} не в LD_LIBRARY_PATH\n"
f" Текущий LD_LIBRARY_PATH: {ld_path or '(не задан)'}\n"
f" Решение: export LD_LIBRARY_PATH={self.native_path}:$LD_LIBRARY_PATH\n"
f" Или в SparkSession: .config('spark.executorEnv.LD_LIBRARY_PATH', '{self.native_path}')"
)
return False
print(f"✅ LD_LIBRARY_PATH содержит путь к нативным библиотекам")
return True
def run_all_checks(self) -> bool:
"""Запускает все проверки и выводит итоговый отчёт."""
print("\n" + "=" * 60)
print("ДИАГНОСТИКА SHORT-CIRCUIT READS")
print("=" * 60)
checks = [
("Наличие Unix Domain Socket", self.check_socket_exists),
("Права доступа на директорию", self.check_socket_permissions),
("Подключение к сокету", self.check_socket_connectable),
("Нативная библиотека libhadoop.so", self.check_libhadoop),
("LD_LIBRARY_PATH", self.check_ld_library_path),
]
all_ok = True
for check_name, check_fn in checks:
print(f"\nПроверка: {check_name}")
try:
result = check_fn()
if not result:
all_ok = False
except Exception as e:
print(f" ⚠️ Ошибка при проверке: {e}")
all_ok = False
if self.warnings:
print("\n" + "-" * 60)
print("ПРЕДУПРЕЖДЕНИЯ:")
for w in self.warnings:
print(f"\n{w}")
if self.issues:
print("\n" + "-" * 60)
print("КРИТИЧЕСКИЕ ПРОБЛЕМЫ:")
for issue in self.issues:
print(f"\n{issue}")
else:
print("\n" + "=" * 60)
print("✅ Все проверки пройдены. Short-Circuit должен работать!")
return all_ok
# ── Запуск диагностики ────────────────────────────────────────────
if __name__ == "__main__":
diag = ShortCircuitDiagnostics(
socket_path="/var/lib/hadoop-hdfs/dn_socket",
hadoop_native_path="/usr/lib/hadoop/lib/native",
)
diag.run_all_checks()
Таблица типичных ошибок и решений¶
| Симптом в логах | Причина | Решение |
|---|---|---|
The domain socket path is insecure |
Права директории слишком широкие (o+w) | chmod o-w /var/lib/hadoop-hdfs |
Unable to load native-hadoop library |
libhadoop.so не найдена или несовместима | Установить hadoop-native или пересобрать с -Pnative |
ConnectionRefusedError на UDS |
DataNode не слушает на сокете | Проверить dfs.client.read.shortcircuit=true в hdfs-site.xml и перезапустить DN |
PermissionError на UDS connect |
Пользователь Spark не в группе hadoop | usermod -aG hadoop spark |
| SC rate = 0% при включённой конфигурации | Spark запущен не на YARN или нет NODE_LOCAL задач | Убедиться что Executor'ы запущены на тех же серверах что и DataNode |
ShortCircuitReplica has expired |
FD Cache слишком мал | Увеличить streams.cache.size до 2000+ |
| Block read через TCP несмотря на SC=true | libhadoop.so не загружена в Executor JVM | Добавить -Djava.library.path в executor.extraJavaOptions |
Итоги: что важно помнить о Short-Circuit Reads¶
Short-Circuit Reads - это одна из тех оптимизаций, которые при правильной настройке дают заметный прирост производительности «бесплатно», без изменения бизнес-логики пайплайна.
Физическая суть: устранение двойного копирования данных через TCP loopback и Context Switch overhead при локальном чтении. Данные идут напрямую: disk → Page Cache → Executor Heap, минуя DataNode heap и TCP буферы.
Три компонента: Unix Domain Socket (для File Descriptor Passing), libhadoop.so (JNI для POSIX системных вызовов), Shared Memory Segments (для zero-roundtrip проверки состояния блоков). Все три должны работать вместе.
Главная ловушка: Short-Circuit никогда не выдаёт ошибку при проблемах - всегда тихий fallback на TCP. Обязательна периодическая проверка JMX-метрик DataNode (ShortCircuitLocalReadsRead / BlocksRead).
Критичные права доступа: директория с Unix Domain Socket должна иметь права без write для «других» (chmod o-w). 777, 775 или 755 с write для other → Short-Circuit не будет работать.
Ограничения применимости: Short-Circuit работает только в HDFS с совмещёнными Compute и Storage. При работе с S3, MinIO, Ceph, Ozone - механизм неприменим. В современных Cloud-Native Lakehouse архитектурах значимость Short-Circuit снижается, но для on-premise HDFS кластеров это по-прежнему важная оптимизация.