메모리 관리
메모리 관리
카프카 스트림즈는 레코드가 상태 저장소에 쓰이거나 다운스트림으로 전달되기 전에 내부 캐싱·컴팩션을 하면서 메모리(RAM)를 사용해요. 이 페이지에서는 DSL과 Processor API에서 레코드 캐시가 각각 어떻게 동작하는지, cache.max.bytes.buffering 설정, 그리고 RocksDB 오프힙(off-heap) 메모리까지 정리해드릴게요.
출처: 문서
본문
레코드의 내부 캐싱과 컴팩션에 사용되는 총 메모리(RAM) 크기를 지정할 수 있어요. 이 캐싱은 레코드가 상태 저장소에 쓰이거나 다른 노드로 다운스트림 전달되기 전에 발생해요.
레코드 캐시는 DSL과 Processor API에서 조금 다르게 구현돼요.
DSL에서의 레코드 캐시
처리 토폴로지 인스턴스의 레코드 캐시에 대한 총 메모리(RAM) 크기를 지정할 수 있어요. 이것은 다음 KTable 인스턴스에 의해 활용돼요:
- 소스
KTable:StreamsBuilder#table()또는StreamsBuilder#globalTable()로 생성된KTable인스턴스. - 집계
KTable: 집계의 결과로 생성된KTable인스턴스.
이러한 KTable 인스턴스에 대해 레코드 캐시는 다음에 사용돼요:
- 기본 상태 저장 프로세서 노드가 내부 상태 저장소에 쓰기 전에 출력 레코드를 내부 캐싱·컴팩션.
- 기본 상태 저장 프로세서 노드가 다운스트림 프로세서 노드로 전달하기 전에 출력 레코드를 내부 캐싱·컴팩션.
다음 예시로 레코드 캐싱이 있을 때와 없을 때의 동작을 이해해보세요. 이 예시에서 입력은 레코드 <K,V>: <A, 1>, <D, 5>, <A, 20>, <A, 300>를 가진 KStream<String, Integer>예요. 이 예시의 초점은 key == A인 레코드에 있어요. 집계가 입력에 대해 키별로 그룹화된 레코드 값의 합을 계산하고 KTable<String, Integer>을 반환해요.
- 캐싱 없음: 결과 집계 테이블의 변화를 나타내는 키
A에 대한 출력 레코드 시퀀스가 생성돼요. 괄호(())는 변화를 나타내고, 왼쪽 숫자는 새 집계 값, 오른쪽 숫자는 이전 집계 값이에요:<A, (1, null)>, <A, (21, 1)>, <A, (321, 21)>. - 캐싱 있음: 키
A에 대해 단일 출력 레코드가 생성되며, 이는 캐시에서 컴팩션될 가능성이 높아<A, (321, null)>의 단일 출력 레코드를 만들어요. 이 레코드는 집계의 내부 상태 저장소에 쓰이고 다운스트림 연산으로 전달돼요.
캐시 크기는 처리 토폴로지당 전역 설정인 cache.max.bytes.buffering 파라미터로 지정돼요:
// Enable record cache of size 10 MB.
Properties props = new Properties();
props.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, 10 * 1024 * 1024L);
이 파라미터는 캐싱에 할당된 바이트 수를 제어해요. 구체적으로, T 스레드와 캐싱에 할당된 C 바이트를 가진 프로세서 토폴로지 인스턴스에 대해, 각 스레드는 균등한 C/T 바이트를 가져 자신의 캐시를 구성하고 태스크 사이에서 적절히 사용해요. 이는 스레드 수만큼의 캐시가 존재하지만 스레드 간 캐시 공유는 없다는 뜻이에요.
캐시의 기본 API는 put()과 get() 호출로 구성돼요. 캐시 크기에 도달한 후 레코드는 간단한 LRU 방식으로 제거돼요. 키가 있는 레코드 R1 = <K1, V1>이 노드에서 처리를 처음 마치면 캐시에서 dirty로 표시돼요. 그 시간 동안 같은 노드에서 처리되는 같은 키 K1을 가진 다른 키가 있는 레코드 R2 = <K1, V2>는 <K1, V1>을 덮어쓸 거예요. 이것을 "컴팩션됨"이라고 해요. 이것은 카프카의 로그 컴팩션과 같은 효과지만, 서버 측(즉 카프카 브로커)이 아니라 클라이언트 측 애플리케이션에서, 레코드가 아직 메모리에 있는 동안 더 일찍 발생해요. 플러시 후 R2는 다음 처리 노드로 전달되고 로컬 상태 저장소에 쓰여요.
캐싱의 의미론은 commit.interval.ms 또는 cache.max.bytes.buffering(캐시 압력) 중 더 이른 것이 도달할 때마다 데이터가 상태 저장소로 플러시되고 다음 다운스트림 프로세서 노드로 전달된다는 거예요. commit.interval.ms와 cache.max.bytes.buffering은 모두 전역 파라미터예요. 그렇기 때문에 개별 노드에 대해 다른 파라미터를 지정할 수 없어요.
원하는 시나리오에 따른 두 파라미터의 예시 설정은 다음과 같아요.
캐싱을 끄려면 캐시 크기를 0으로 설정할 수 있어요:
// Disable record cache
Properties props = new Properties();
props.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, 0);
캐싱을 활성화하되 레코드가 캐시되는 기간에 상한을 두고 싶다면 커밋 간격을 설정할 수 있어요. 이 예시에서는 1000밀리초로 설정했어요:
Properties props = new Properties();
// Enable record cache of size 10 MB.
props.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, 10 * 1024 * 1024L);
// Set commit interval to 1 second.
props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000);
이 두 구성의 효과는 아래 그림에 설명돼 있어요. 레코드는 4개의 키(파랑, 빨강, 노랑, 초록)로 표시돼요. 캐시에 3개의 키만을 위한 공간이 있다고 가정해요.
- 캐시가 비활성화되면 (a), 모든 입력 레코드가 출력돼요.
- 캐시가 활성화되면 (b):
- 대부분의 레코드는 커밋 간격이 끝날 때 출력돼요 (예:
t1에서 파란 레코드 하나가 출력되는데, 이는 그 시점까지 파란 키의 최종 덮어쓰기예요). - 일부 레코드는 캐시 압력 때문에 출력돼요 (즉 커밋 간격이 끝나기 전에). 예를 들어
t2전의 빨간 레코드를 보세요. 더 작은 캐시 크기에서는 레코드가 언제 출력되는지를 결정하는 주 요인이 캐시 압력일 것으로 예상해요. 큰 캐시 크기에서는 커밋 간격이 주 요인이 될 거예요. - 출력된 레코드의 총 수가 15개에서 8개로 줄었어요.
- 대부분의 레코드는 커밋 간격이 끝날 때 출력돼요 (예:
Processor API에서의 레코드 캐시
처리 토폴로지 인스턴스의 레코드 캐시에 대한 총 메모리(RAM) 크기를 지정할 수 있어요. 이것은 상태 저장 프로세서 노드에서 상태 저장소로 쓰기 전에 출력 레코드를 내부 캐싱·컴팩션하는 데 사용돼요.
Processor API의 레코드 캐시는 다운스트림으로 전달되는 출력 레코드를 캐싱하거나 컴팩션하지 않아요. 즉 모든 다운스트림 프로세서 노드는 모든 레코드를 볼 수 있는 반면, 상태 저장소는 줄어든 수의 레코드를 봐요. 이것은 시스템의 정확성에 영향을 주지 않지만, 상태 저장소를 위한 성능 최적화예요. 예를 들어, Processor API로 다른 값을 다운스트림으로 전달하면서 레코드를 상태 저장소에 저장할 수 있어요.
State Stores 섹션에 처음 나온 예시에 이어서, 캐싱을 비활성화하려면 withCachingDisabled 호출을 추가할 수 있어요 (캐시는 기본으로 활성화되어 있지만 명시적인 withCachingEnabled 호출도 있어요):
StoreBuilder countStoreBuilder =
Stores.keyValueStoreBuilder(
Stores.persistentKeyValueStore("Counts"),
Serdes.String(),
Serdes.Long())
.withCachingEnabled();
레코드 캐시는 버전 상태 저장소(versioned state store)에서는 지원되지 않아요. 오래된 데이터를 읽지 않으려면 반복자(iterator)를 만들기 전에 저장소를 flush()할 수 있어요. 너무 자주 플러시하면 RocksDB를 사용할 때 성능 저하가 발생할 수 있으므로, 일반적으로 수동 플러시는 피하는 것이 좋아요.
RocksDB
각 RocksDB 인스턴스는 블록 캐시, 인덱스·필터 블록, memtable(쓰기 버퍼)에 오프힙(off-heap) 메모리를 할당해요. 중요한 구성(RocksDB 4.1.0 기준)은 block_cache_size, write_buffer_size, max_write_buffer_number를 포함해요. 이것들은 rocksdb.config.setter 구성을 통해 지정할 수 있어요.
체인지로그 오프셋 내구성과 플러시 빈도
4.3(KIP-1035)부터 Kafka Streams는 각 영속 저장소의 체인지로그 오프셋을 태스크별 .checkpoint 파일 대신 RocksDB 내부에 저장해요. Kafka Streams는 WAL(write-ahead log)을 비활성화한 상태로 RocksDB를 실행해요 — 체인지로그 토픽이 레코드의 내구성 로그이므로 — 데이터와 체인지로그 오프셋은 memtable이 SST 파일로 플러시될 때에만 디스크에서 내구성이 생겨요. 플러시는 memtable이 write_buffer_size(기본 16 MB)를 채우거나 저장소가 정상적으로 닫힐 때(KafkaStreams#close) 발생해요. 이전 릴리스와 달리 Kafka Streams는 더 이상 매 커밋마다 RocksDB를 강제 플러시하지 않아요.
실질적인 의미는 디스크의 체인지로그 오프셋은 마지막 유기적 플러시 또는 정상 종료 때만큼만 최신이라는 거예요. 고처리량 저장소의 경우 memtable이 자주 채워지므로 이것은 문제가 되지 않아요. memtable을 채우는 데 오래 걸릴 수 있는 저트래픽 저장소의 경우 영속된 오프셋이 저장소의 실제 위치보다 오랫동안 뒤처질 수 있어요. 애플리케이션이 비정상적으로 종료되면(예: SIGKILL, OOM-킬, 또는 프로세스/포드 종료 유예 기간 내에 완료되지 않는 KafkaStreams#close), 마지막으로 플러시된 오프셋만 남아요. 그 사이 체인지로그 토픽의 로그 시작 오프셋이 그 오래된 오프셋을 넘어 진행되었다면(리텐션 또는 컴팩션을 통해), 복원 컨슈머가 다음 재시작 시 범위를 벗어난 위치를 찾게 되고, 태스크는 체인지로그에서 재초기화돼요(OffsetOutOfRangeException / TaskCorruptedException으로 기록). 이것은 데이터 손실 없이 자동으로 복구되지만, 태스크의 전체 재-복원을 초래해요.
영향을 받는 저트래픽 저장소에 대해 이 가능성을 줄이려면, 커스텀 RocksDBConfigSetter를 통해 write_buffer_size를 줄여(그리고 max_write_buffer_number도 조정) 더 자주 플러시하게 할 수 있어요:
public static class CustomRocksDBConfig implements RocksDBConfigSetter {
@Override
public void setConfig(final String storeName, final Options options, final Map<String, Object> configs) {
// smaller write buffer => more frequent flushes => fresher on-disk changelog offset,
// at the cost of more (smaller) SST files and more compaction work
options.setWriteBufferSize(4 * 1024 * 1024L);
}
@Override
public void close(final String storeName, final Options options) {}
}
이것은 트레이드오프예요: 더 작은 쓰기 버퍼는 영속된 오프셋이 얼마나 오래될 수 있는지를 제한하지만, SST 파일 수와 컴팩션 부하를 증가시켜요. 그러니 영향을 받는 특정 저트래픽 저장소만 튜닝하고, 버퍼를 저장소의 쓰기 속도에 맞게 크기를 정하세요. 플러시는 시간이 아닌 용량(volume) 기반이에요. 약간의 쓰기만 받는 저장소는 닫힐 때까지 플러시되지 않을 수도 있으므로, 가장 신뢰할 수 있는 보호는 정상 종료예요. KafkaStreams#close를 호출하는 셧다운 훅을 등록하고 완료되도록 허용하세요. 쿠버네티스에서는 종료 타임아웃을 포드 종료 유예 기간(기본 30초)보다 충분히 아래로 설정해 프로세스가 종료 중 SIGKILL되지 않게 하세요. 저장소 앞의 레코드 캐시(statestore.cache.max.bytes)는 RocksDB에 쓰기 전에 같은 키에 대한 반복 업데이트를 제자리에서 합쳐(coalesce)요. 그래서 업데이트가 많고 키 공간이 작은 워크로드에서는 커밋당 memtable에 도달하는 바이트가 더 적어 플러시가 더 느려져요.
또한 RocksDB의 기본 메모리 할당자를 바꾸는 것을 권장해요. 기본 할당자는 메모리 소비를 증가시킬 수 있기 때문이에요. 메모리 할당자를 jemalloc으로 바꾸려면 Kafka Streams 애플리케이션을 시작하기 전에 환경 변수 LD_PRELOAD를 설정해야 해요:
# example: install jemalloc (on Debian)
$ apt install -y libjemalloc-dev
# set LD_PRELOAD before you start your Kafka Streams application
$ export LD_PRELOAD="/usr/lib/x86_64-linux-gnu/libjemalloc.so"
2.3.0부터 모든 인스턴스에 걸친 메모리 사용량을 제한해 Kafka Streams 애플리케이션의 총 오프힙 메모리를 묶을 수 있어요. 그러려면 RocksDB가 인덱스·필터 블록을 블록 캐시에 캐시하도록 구성하고, 공유 WriteBufferManager를 통해 memtable 메모리를 제한하고 그 메모리를 블록 캐시에 계산하며, 같은 Cache 객체를 각 인스턴스에 전달해야 해요. 자세한 내용은 RocksDB Memory Usage를 참고해요. 이것을 구현하는 예시 RocksDBConfigSetter는 다음과 같아요:
public static class BoundedMemoryRocksDBConfig implements RocksDBConfigSetter {
private static org.rocksdb.Cache cache = new org.rocksdb.LRUCache(TOTAL_OFF_HEAP_MEMORY, -1, false, INDEX_FILTER_BLOCK_RATIO);1
private static org.rocksdb.WriteBufferManager writeBufferManager = new org.rocksdb.WriteBufferManager(TOTAL_MEMTABLE_MEMORY, cache);
@Override
public void setConfig(final String storeName, final Options options, final Map<String, Object> configs) {
BlockBasedTableConfig tableConfig = (BlockBasedTableConfig) options.tableFormatConfig();
// These three options in combination will limit the memory used by RocksDB to the size passed to the block cache (TOTAL_OFF_HEAP_MEMORY)
tableConfig.setBlockCache(cache);
tableConfig.setCacheIndexAndFilterBlocks(true);
options.setWriteBufferManager(writeBufferManager);
// These options are recommended to be set when bounding the total memory
tableConfig.setCacheIndexAndFilterBlocksWithHighPriority(true);2
tableConfig.setPinTopLevelIndexAndFilter(true);
tableConfig.setBlockSize(BLOCK_SIZE);3
options.setMaxWriteBufferNumber(N_MEMTABLES);
options.setWriteBufferSize(MEMTABLE_SIZE);
options.setTableFormatConfig(tableConfig);
}
@Override
public void close(final String storeName, final Options options) {
// Cache and WriteBufferManager should not be closed here, as the same objects are shared by every store instance.
}
}
INDEX_FILTER_BLOCK_RATIO는 블록 캐시의 일부를 "high priority"(즉 인덱스와 필터) 블록용으로 떼어 놓아 데이터 블록에 의해 제거되지 않도록 하는 데 사용할 수 있어요. 캐시 생성자의 boolean 파라미터는 용량보다 커질 수 있는 드문 경우에 캐시가 읽기·반복을 실패시킴으로써 엄격한 메모리 제한을 강제할지 제어해요.INDEX_FILTER_BLOCK_RATIO가 효과를 가지려면(각주 1 참조) 이것을 설정해야 해요.- RocksDB 문서의 지침에 따라 기본 블록 크기를 수정하고 싶을 수 있어요. 더 큰 블록 크기는 인덱스 블록이 더 작아지지만, 캐시된 데이터 블록은 그렇지 않으면 제거되었을 더 많은 콜드 데이터를 포함할 수 있어요.
참고: 적어도 위 구성을 설정하는 것을 권장하지만, 최상의 성능을 내는 특정 옵션은 워크로드에 따라 달라져요. 특정 사용 사례에 가장 좋은 선택을 정하기 위해 실험해보는 것이 좋아요. 한 앱에 최적인 구성이 다른 토폴로지나 입력 토픽을 가진 앱에는 적용되지 않을 수 있다는 점을 명심하세요. 위 권장 구성 외에 RocksDB 문서에 설명된 대로 partitioned index filter를 사용하는 것도 고려할 수 있어요.
기타 메모리 사용량
Apache Kafka 내부에 런타임 중 메모리를 할당하는 다른 모듈이 있어요. 다음을 포함해요:
- 프로듀서 버퍼링: 프로듀서 구성
buffer.memory로 관리. - 컨슈머 버퍼링: 현재 엄격히 관리되지 않지만, 페치 크기, 즉
fetch.max.bytes와fetch.max.wait.ms로 간접적으로 제어 가능. - 프로듀서와 컨슈머 모두 버퍼링 메모리로 계산되지 않는 별도의 TCP 송신/수신 버퍼도 있어요. 이것들은
send.buffer.bytes/receive.buffer.bytes구성으로 제어돼요. - 역직렬화된 객체 버퍼링:
consumer.poll()이 레코드를 반환한 후, 그것들은 타임스탬프를 추출하기 위해 역직렬화되고 스트림즈 공간에 버퍼링돼요. 현재 이것은buffered.records.per.partition으로만 간접적으로 제어돼요.
팁: 반복자(iterator)는 리소스를 해제하기 위해 명시적으로 닫아야 해요. 저장소 반복자(예:
KeyValueIterator,WindowStoreIterator)는 완료 시 파일 핸들러와 인메모리 읽기 버퍼 같은 리소스를 해제하기 위해 명시적으로 닫아야 하며, 이 Closeable 클래스에 try-with-resources 문(JDK7부터 사용 가능)을 사용해도 돼요. 그렇지 않으면 스트림 애플리케이션의 메모리 사용량이 OOM에 도달할 때까지 실행 중 계속 증가해요.
더 알아보기
- Streams 애플리케이션 구성하기 —
cache.max.bytes.buffering,commit.interval.ms등 전체 구성을 봐요. - 인터랙티브 쿼리 — 상태 저장소를 쿼리하는 방법을 봐요.
- Streams 애플리케이션 실행하기 — 운영 시 메모리 관련 고려사항을 봐요.