Spark 메모리 관리와 HDFS의 운영체제 상호작용
Spark의 통합 메모리와 스필링, HDFS의 파일 시스템·페이지 캐시·데이터 지역성 상호작용을 정리한다.
2026-08-14 · 최초 발행 2026-01-19
JVM 힙 안에서 실행과 캐시가 만나는 지점
Spark의 인메모리 처리는 빠르지만, 실제 처리량은 JVM 힙을 실행 작업과 캐시 데이터가 어떻게 나눠 쓰는지에 따라 달라진다. 통합 메모리 모델에서는 실행 메모리와 저장 메모리가 고정된 벽으로 분리되지 않고 같은 영역에서 동적으로 사용된다.
실행 메모리는 셔플·조인·정렬·집계 같은 연산에 쓰인다. 저장 메모리는 RDD 캐싱과 브로드캐스트 변수를 담는다. 이 밖에 사용자 정의 데이터 구조와 메타데이터를 위한 User Memory, Spark 내부 객체를 위한 Reserved Memory가 있다.
다음 설정은 executor와 driver의 메모리, 메모리 영역 비율, 셔플 파티션 수를 함께 지정하고 executor 저장소 상태를 출력한다.
# PySpark 메모리 설정 예시
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("MemoryOptimization") \
.config("spark.executor.memory", "4g") \
.config("spark.executor.memoryOverhead", "512m") \
.config("spark.driver.memory", "2g") \
.config("spark.memory.fraction", "0.6") \
.config("spark.memory.storageFraction", "0.5") \
.config("spark.sql.shuffle.partitions", "200") \
.getOrCreate()
# 메모리 사용 모니터링
def print_memory_stats():
"""Executor 메모리 통계 출력"""
storage_status = spark.sparkContext._jsc.sc().getExecutorStorageStatus()
for status in storage_status:
print(f"Executor: {status.blockManagerId().executorId()}")
print(f" Max Memory: {status.maxMem() / 1024 / 1024:.2f} MB")
print(f" Used Memory: {status.memUsed() / 1024 / 1024:.2f} MB")
print(f" Remaining: {status.memRemaining() / 1024 / 1024:.2f} MB")
캐시 방식은 데이터 재사용 빈도와 메모리 여유에 맞춰 골라야 한다. 메모리에만 보관할지, 부족할 때 디스크로 내보낼지, 직렬화나 off-heap을 활용할지를 StorageLevel로 정한다.
# 다양한 캐싱 레벨
from pyspark import StorageLevel
# 데이터 로드
data = spark.read.parquet("hdfs://namenode:9000/large_dataset")
# MEMORY_ONLY: 메모리에만 저장 (직렬화 안함)
data.persist(StorageLevel.MEMORY_ONLY)
# MEMORY_AND_DISK: 메모리 부족 시 디스크로 스필
data.persist(StorageLevel.MEMORY_AND_DISK)
# MEMORY_ONLY_SER: 직렬화하여 메모리에 저장 (공간 효율적)
data.persist(StorageLevel.MEMORY_ONLY_SER)
# DISK_ONLY: 디스크에만 저장
data.persist(StorageLevel.DISK_ONLY)
# OFF_HEAP: JVM 힙 외부에 저장 (GC 부담 감소)
spark.conf.set("spark.memory.offHeap.enabled", "true")
spark.conf.set("spark.memory.offHeap.size", "2g")
data.persist(StorageLevel.OFF_HEAP)
# 캐시 해제
data.unpersist()
GC 설정과 로그는 executor 메모리 문제가 JVM 수준의 회수 지연과 관련 있는지 판단할 때 함께 본다.
# G1GC 설정 (Spark 권장)
spark-submit \
--conf spark.executor.extraJavaOptions="-XX:+UseG1GC -XX:InitiatingHeapOccupancyPercent=35 -XX:ConcGCThreads=12" \
--conf spark.driver.extraJavaOptions="-XX:+UseG1GC" \
my_spark_job.py
# GC 로그 활성화
spark-submit \
--conf spark.executor.extraJavaOptions="-XX:+PrintGCDetails -XX:+PrintGCTimeStamps -XX:+PrintGCDateStamps -Xloggc:/var/log/spark/gc.log" \
my_spark_job.py
메모리 부족이 디스크 작업으로 바뀌는 과정
연산 중 메모리 버퍼가 차면 Spark는 데이터를 정렬해 spill file로 기록하고 버퍼를 비운다. 이후 메모리에서 처리된 데이터와 spill 파일을 최종 병합한다. 이 전환은 셔플과 집계 작업에서 디스크 I/O로 이어진다.
스필링은 Spark UI의 stage 메트릭으로 확인할 수 있다. 아래 코드는 job에 연결된 stage별 메모리 및 디스크 스필 바이트를 출력하고, 이어서 셔플 파티션과 적응형 쿼리 실행을 설정한다.
# 스필링 발생 확인
def check_spilling(spark_job_id):
"""Spark UI에서 스필링 메트릭 확인"""
# 스파크 메트릭 접근 (실제로는 Spark UI REST API 사용)
stages = spark.sparkContext.statusTracker().getJobInfo(spark_job_id).stageIds()
for stage_id in stages:
stage_info = spark.sparkContext.statusTracker().getStageInfo(stage_id)
if stage_info:
# 메모리 스필 정보
print(f"Stage {stage_id}:")
print(f" Memory Bytes Spilled: {stage_info.memoryBytesSpilled}")
print(f" Disk Bytes Spilled: {stage_info.diskBytesSpilled}")
# 스필링 최소화 전략
spark.conf.set("spark.sql.shuffle.partitions", "400") # 파티션 수 증가
spark.conf.set("spark.sql.adaptive.enabled", "true") # 적응형 쿼리 실행
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
HDFS는 OS 파일 시스템 위에 블록을 남긴다
HDFS에서 NameNode는 네임스페이스 이미지, Edit Log, 블록 매핑 같은 메타데이터를 담당하고, DataNode는 실제 블록 복제본을 보관한다. 클라이언트는 메타데이터 조회와 데이터 전송에서 이 역할을 구분해 통신한다.
HDFS 기본 블록 크기는 128MB이고 OS 파일 시스템은 4KB 단위로 동작한다. HDFS 블록은 OS의 일반 파일로 저장되며, NameNode는 메타데이터를 메모리에 유지한다. 블록 위치는 작업 스케줄링의 데이터 지역성에 반영되고, 기본 3개의 복제본은 내결함성을 제공한다.
읽기 작업에서는 NameNode로부터 블록 위치를 얻은 뒤 DataNode에서 데이터를 직접 읽는다. 쓰기 시에는 복제 수와 블록 크기를 지정할 수 있다.
# HDFS Python 클라이언트 (hdfs3 라이브러리 사용)
from hdfs3 import HDFileSystem
hdfs = HDFileSystem(host='namenode', port=9000)
def read_hdfs_file(path: str):
"""HDFS 파일 읽기"""
# 1. NameNode에서 블록 위치 정보 가져오기
file_info = hdfs.info(path)
print(f"File size: {file_info['size']} bytes")
# 2. 블록별 위치 정보
block_locations = hdfs.get_block_locations(path)
for block in block_locations:
print(f"Block offset: {block['offset']}")
print(f"Block length: {block['length']}")
print(f"DataNodes: {block['hosts']}")
# 3. 데이터 읽기 (DataNode로부터 직접)
with hdfs.open(path, 'rb') as f:
data = f.read()
return data
def write_hdfs_file(path: str, data: bytes):
"""HDFS 파일 쓰기"""
with hdfs.open(path, 'wb', replication=3, block_size=134217728) as f:
f.write(data)
쓰기 파이프라인에서는 클라이언트가 NameNode에서 블록을 할당받고 DataNode 간에 데이터를 전달한다. ACK는 마지막 DataNode에서 앞선 노드와 클라이언트 방향으로 돌아간다.
페이지 캐시와 로컬 읽기 경로
DataNode가 자주 읽는 블록은 OS 페이지 캐시에 남을 수 있다. HDFS의 메모리 동작을 볼 때 JVM 프로세스만이 아니라 운영체제가 확보한 buff/cache도 함께 확인해야 하는 이유다. 아래 설정은 페이지 캐시 확인과 테스트 목적의 캐시 삭제, DataNode locked memory 설정 예시를 담고 있다.
# 페이지 캐시 확인
free -h
# total used free shared buff/cache available
# Mem: 62Gi 15Gi 2.0Gi 100Mi 45Gi 46Gi
# HDFS는 OS 페이지 캐시를 적극 활용
# DataNode가 자주 읽는 블록은 자동으로 캐시됨
# 페이지 캐시 삭제 (테스트 목적)
echo 3 > /proc/sys/vm/drop_caches
# HDFS DataNode 메모리 설정
# hdfs-site.xml
# <property>
# <name>dfs.datanode.max.locked.memory</name>
# <value>8589934592</value> <!-- 8GB -->
# </property>
같은 노드에 데이터가 있을 때 short-circuit local read를 사용하면 네트워크 경로 대신 로컬 읽기 경로를 활용할 수 있다.
<!-- hdfs-site.xml: 로컬 읽기 최적화 -->
<configuration>
<property>
<name>dfs.client.read.shortcircuit</name>
<value>true</value>
</property>
<property>
<name>dfs.domain.socket.path</name>
<value>/var/lib/hadoop-hdfs/dn_socket</value>
</property>
<property>
<name>dfs.client.read.shortcircuit.skip.checksum</name>
<value>false</value>
</property>
</configuration>
DataNode 프로세스의 네트워크 인터페이스 통계는 HDFS 전송량을 점검하는 한 방법이다.
# HDFS 네트워크 성능 모니터링
import subprocess
import re
def monitor_hdfs_network():
"""DataNode 네트워크 사용량 모니터링"""
# DataNode 프로세스 찾기
pid = subprocess.check_output(
"pgrep -f 'DataNode'",
shell=True
).decode().strip()
# 네트워크 통계
result = subprocess.check_output(
f"cat /proc/{pid}/net/dev",
shell=True
).decode()
# 수신/송신 바이트 파싱
for line in result.split('\n'):
if 'eth0' in line or 'ens' in line:
parts = line.split()
rx_bytes = int(parts[1])
tx_bytes = int(parts[9])
print(f"RX: {rx_bytes / 1024 / 1024:.2f} MB")
print(f"TX: {tx_bytes / 1024 / 1024:.2f} MB")
Spark 스케줄러가 HDFS 블록 위치를 이용하는 방식
Spark Driver는 NameNode에서 블록 위치 정보를 받은 뒤 Task Scheduler가 지역성 수준을 고려해 작업을 배치한다. 같은 JVM, 같은 노드, 같은 랙, 다른 랙 순서로 데이터와 가까운 실행 위치를 판단한다.
다음 코드는 HDFS의 Parquet 데이터를 읽고 지역성 대기 시간을 설정한 후 집계 작업을 실행한다.
# Spark에서 데이터 지역성 확인
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("LocalityTest").getOrCreate()
# HDFS에서 데이터 읽기
data = spark.read.parquet("hdfs://namenode:9000/data")
# 데이터 지역성 설정
spark.conf.set("spark.locality.wait", "3s") # 지역성 대기 시간
spark.conf.set("spark.locality.wait.node", "5s")
spark.conf.set("spark.locality.wait.rack", "10s")
# 작업 실행 및 메트릭 확인
result = data.groupBy("category").count()
result.show()
# Spark UI에서 Locality Level 확인
# http://spark-driver:4040/stages/
HDFS 읽기 설정, 파일 파티션 크기, Parquet 옵션도 Spark 작업의 입력 경로에 직접 영향을 준다.
# Spark에서 HDFS 읽기 최적화
spark.conf.set("spark.hadoop.dfs.client.read.shortcircuit", "true")
spark.conf.set("spark.hadoop.dfs.client.read.shortcircuit.skip.checksum", "false")
spark.conf.set("spark.hadoop.dfs.client.cache.drop.behind.reads", "true")
# 파티션 크기 조정
spark.conf.set("spark.sql.files.maxPartitionBytes", "134217728") # 128MB
spark.conf.set("spark.sql.files.openCostInBytes", "4194304") # 4MB
# Parquet 읽기 최적화
spark.conf.set("spark.sql.parquet.compression.codec", "snappy")
spark.conf.set("spark.sql.parquet.filterPushdown", "true")
spark.conf.set("spark.sql.parquet.mergeSchema", "false")
메모리·디스크·커널 설정을 함께 점검할 때
off-heap 메모리는 JVM 힙 외부를 사용하는 선택지다. 다음 실행 옵션과 설정은 off-heap, executor memory overhead, Project Tungsten 관련 항목을 보여 준다.
# Off-Heap 메모리 설정
spark-submit \
--conf spark.memory.offHeap.enabled=true \
--conf spark.memory.offHeap.size=4g \
--conf spark.executor.memoryOverhead=2g \
my_job.py
# Project Tungsten: 바이트코드 생성 및 off-heap 최적화
# 자동으로 활성화되지만 명시적 설정 가능
spark.conf.set("spark.sql.codegen.wholeStage", "true")
spark.conf.set("spark.sql.tungsten.enabled", "true")
DataNode의 데이터 디렉터리를 여러 디스크에 두고 전송 스레드와 블록 보고서 관련 값을 조정하는 설정은 다음과 같다.
<!-- hdfs-site.xml: 디스크 I/O 최적화 -->
<configuration>
<!-- 여러 디스크 사용 -->
<property>
<name>dfs.datanode.data.dir</name>
<value>/data1/hdfs,/data2/hdfs,/data3/hdfs,/data4/hdfs</value>
</property>
<!-- 동시 전송 스레드 수 -->
<property>
<name>dfs.datanode.max.transfer.threads</name>
<value>8192</value>
</property>
<!-- 블록 보고서 전송 최적화 -->
<property>
<name>dfs.blockreport.split.threshold</name>
<value>1000000</value>
</property>
</configuration>
디스크 스케줄러, 파일 디스크립터 한도, 스왑, 투명 대형 페이지, TCP 버퍼 역시 Spark와 Hadoop 워크로드가 놓인 OS 환경의 변수다.
# I/O 스케줄러 설정 (SSD용 noop, HDD용 deadline)
echo noop > /sys/block/sda/queue/scheduler
# 파일 디스크립터 한도 증가
ulimit -n 65536
# 스왑 최소화
sysctl -w vm.swappiness=1
# 투명 대형 페이지(THP) 비활성화 (Hadoop/Spark 권장)
echo never > /sys/kernel/mm/transparent_hugepage/enabled
echo never > /sys/kernel/mm/transparent_hugepage/defrag
# TCP 튜닝 (네트워크 처리량 개선)
sysctl -w net.core.rmem_max=134217728
sysctl -w net.core.wmem_max=134217728
sysctl -w net.ipv4.tcp_rmem="4096 87380 67108864"
sysctl -w net.ipv4.tcp_wmem="4096 65536 67108864"
메모리 부족과 작은 파일을 다루는 코드
OutOfMemoryError가 발생한 경우에는 파티션 수, 브로드캐스트 조인 임계값, 읽을 컬럼, 데이터 타입, 중간 캐시, 배치 처리 방식을 함께 검토할 수 있다.
# 문제: OutOfMemoryError 발생
# 해결 전략:
# 1. 파티션 수 증가 (데이터 분산)
spark.conf.set("spark.sql.shuffle.partitions", "1000")
# 2. 브로드캐스트 조인 임계값 조정
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "10485760") # 10MB
# 3. 불필요한 컬럼 제거
df = spark.read.parquet("hdfs://data").select("id", "name", "value")
# 4. 데이터 타입 최적화
from pyspark.sql.types import *
df = df.withColumn("id", df["id"].cast(IntegerType()))
# 5. 중간 결과 캐싱 제거
df.unpersist()
# 6. 점진적 처리
def process_in_batches(df, batch_size=10000):
total = df.count()
for offset in range(0, total, batch_size):
batch = df.offset(offset).limit(batch_size)
process_batch(batch)
수백만 개의 작은 파일은 NameNode 메모리 부족으로 이어질 수 있다. 작은 파일을 읽어 큰 파일로 다시 기록하거나 Hadoop 네이티브 도구를 사용하는 방법은 다음과 같다.
# 문제: 수백만 개의 작은 파일로 인한 NameNode 메모리 부족
# 해결: 파일 병합
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("FileMerge").getOrCreate()
# 작은 파일들을 읽어서 큰 파일로 재작성
input_path = "hdfs://namenode:9000/small_files/*"
output_path = "hdfs://namenode:9000/merged_files"
df = spark.read.parquet(input_path)
# 파티션 수 조정 (예: 1000개 파티션 → 100개 파티션)
df.coalesce(100).write.mode("overwrite").parquet(output_path)
# 또는 Hadoop 네이티브 도구 사용
# hadoop jar hadoop-streaming.jar -D mapreduce.job.reduces=100 \
# -input /small_files -output /merged_files \
# -mapper cat -reducer cat