Devin.KR

스트림 기초 - 큰 파일을 메모리에 올리지 않고 처리하기

개발자KR 조회 18

이 장에서 배우는 것

기본서에서는 파일 시스템 모듈을 사용해 접속 로그를 디스크에 저장하고 다시 불러오는 방법을 다뤘다. 하지만 서비스 규모가 커져서 하루에 쌓이는 로그 파일이 수 기가바이트(GB) 단위로 늘어나면, 전체 파일을 한 번에 메모리에 불러와 처리하는 기존 방식은 한계에 부딪힌다. Node.js는 아주 큰 파일이나 네트워크 데이터를 다룰 때 메모리를 효율적으로 사용할 수 있도록 데이터를 조각내어 흐르게 하는 스트림(Stream)을 제공한다. 이 장에서는 스트림을 활용해 대용량 로그 파일을 안전하게 분석하는 과정을 알아본다.

  • 한 번에 읽기 방식과 스트림 읽기 방식의 메모리 점유율 차이를 설명한다.
  • 읽기, 쓰기, 변환 스트림의 역할과 데이터가 흘러가는 방향을 이해한다.
  • 파이프라인을 구축하여 여러 스트림을 연결하고 오류를 안전하게 처리한다.
  • 텍스트 파일을 줄 단위로 나누어 메모리 누수 없이 가공한다.

문제 상황

서버에 어제 날짜의 전체 접속 로그 파일이 3GB 크기로 저장되어 있다. 운영팀에서 이 파일 중 HTTP 상태 코드가 500번대인 서버 오류 로그만 분리하여 별도 파일로 만들어 달라고 요청했다. 앞 장에서 배운 fs.readFileSync나 fs.readFile을 사용해 3GB 파일을 문자열 변수에 통째로 할당한 다음 split('\n')으로 줄 단위로 나누는 코드를 작성하면, 실행 즉시 프로세스가 강제 종료된다.

Node.js의 기반인 V8 자바스크립트 엔진은 하나의 프로세스가 사용할 수 있는 힙(Heap) 메모리 한도가 시스템에 따라 보통 2GB에서 4GB 사이로 제한되어 있다. 따라서 3GB 파일을 한 번에 읽으려 시도하면 메모리 할당 실패 오류가 발생한다. 설령 실행 인자를 통해 힙 한도를 강제로 늘리더라도, 로그 분석 작업이 실행되는 동안 서버의 남은 메모리가 고갈되어 다른 사용자 요청을 처리하지 못하는 병목 현상이 발생한다.

메모리 사용량 비교: 전체 읽기와 스트림 읽기

파일을 다루는 방식은 데이터를 메모리에 언제, 얼마나 올릴 것인가에 따라 크게 두 가지로 나뉜다. 첫 번째는 전체 파일 데이터를 단일 버퍼나 문자열로 모두 읽어들인 뒤에 다음 작업을 시작하는 방식이다.

파일 읽기 방식의 동작과 자원 소모 비교
구분 전체 읽기 (fs.readFile) 스트림 읽기 (fs.createReadStream)
메모리 점유율 파일 전체 크기 이상 (문자열 파싱 시 추가 소모) 사전에 설정된 청크 크기(보통 64KB)만 유지
데이터 처리 시점 파일 읽기가 완전히 끝난 후 일정 크기의 청크를 읽어들인 즉시
안정성 파일 크기가 V8 힙 한도를 넘으면 프로세스 종료 파일 크기가 아무리 커져도 메모리가 일정하게 유지

반면 스트림은 물이 파이프를 통과하듯, 데이터를 작은 조각(청크, Chunk)으로 나누어 메모리에 올린다. 64KB 크기의 청크 하나를 읽어들여 처리를 마치면, 가비지 컬렉터가 해당 메모리를 회수하고 다음 64KB 청크를 읽어들인다. 이 방식을 사용하면 수십 기가바이트의 파일도 수십 메가바이트의 메모리만으로 안정하게 분석할 수 있다.

전체 파일 읽기는 힙 메모리를 초과하지만 스트림은 작은 조각으로 나누어 안정적으로 처리한다

스트림의 종류와 데이터 흐름

표준 내장 모듈인 node:stream은 목적에 맞춰 세분화된 세 가지 주요 스트림 클래스를 제공한다. 데이터가 어디서 출발하고 어디로 도착하는지 그 흐름을 파악하는 것이 중요하다.

  • 읽기 스트림(Readable Stream): 데이터의 출발지 역할을 한다. 디스크의 파일을 읽어오거나 HTTP 서버에서 클라이언트의 요청 본문을 수신할 때 사용한다. 파일 시스템의 fs.createReadStream이 대표적이다.
  • 쓰기 스트림(Writable Stream): 데이터의 목적지 역할을 한다. 파일에 처리된 데이터를 저장하거나 클라이언트에게 HTTP 응답 본문을 전송할 때 사용한다. fs.createWriteStream으로 생성한다.
  • 변환 스트림(Transform Stream): 읽기 스트림과 쓰기 스트림 사이에 위치하여 지나가는 데이터를 가공한다. 들어온 청크를 압축, 암호화, 또는 특정 조건에 따라 걸러내는(필터링) 작업을 수행한다.

우리가 해결해야 하는 로그 추출 작업은 원본 로그 파일에서 읽기 스트림을 열고, 500번대 상태 코드를 찾는 변환 스트림을 거친 뒤, 조건에 맞는 데이터만 대상 파일의 쓰기 스트림으로 밀어넣는 흐름을 가진다. 텍스트 파일을 한 줄씩 끊어주는 node:readline 모듈 또한 내부적으로 데이터를 변환하는 역할을 수행한다.

읽기 스트림에서 변환 스트림을 거쳐 쓰기 스트림으로 이어지는 데이터 파이프라인 흐름

파이프라인과 안전한 오류 처리

출발지에서 나온 데이터를 목적지로 자연스럽게 흘려보내는 작업을 파이핑(Piping)이라고 부른다. 초기 Node.js 환경에서는 readStream.pipe(writeStream)과 같은 메서드 체이닝을 주로 사용했다. 하지만 이 방식은 스트림 전송 중간에 네트워크 단절이나 파일 권한 문제로 오류가 발생했을 때, 연결된 다른 스트림 객체가 닫히지 않고 열려 있어 메모리 누수를 일으키는 단점이 있었다.

현재는 node:stream/promises 모듈이 제공하는 pipeline 유틸리티 함수를 사용하는 것을 권장한다. 여러 스트림을 인자로 넘겨주면 내부적으로 데이터를 안전하게 전달하며, 과정 중 단 하나의 스트림이라도 오류를 발생시키면 관련된 모든 스트림 자원을 즉시 회수하고 연결을 종료한다.

오류 전파 방식과 자원 회수 비교
메서드 종류 오류 이벤트 처리 방식 자원 정리 및 파괴 (Destroy)
.pipe() 연결된 모든 스트림마다 error 리스너를 개별 등록해야 함 개발자가 직접 오류 상황을 판단하여 각 스트림을 수동으로 파괴해야 함
pipeline() 마지막에 반환되는 프라미스의 catch 블록 한 곳에서 통합 처리 내부 로직이 오류를 감지하면 연관된 모든 스트림을 자동으로 닫음

완성 코드

다음 스크립트는 가상의 로그 파일을 생성한 뒤, 스트림과 파이프라인을 사용해 상태 코드가 500번대인 로그 줄만 걸러내어 새로운 파일에 기록한다.

extract-errors.mjs

import { createReadStream, createWriteStream } from 'node:fs';
import { pipeline } from 'node:stream/promises';
import { createInterface } from 'node:readline';
import { Transform } from 'node:stream';

const SOURCE_FILE = 'server-access.log';
const TARGET_FILE = 'server-errors.log';

// 로그 파일 형식을 흉내 낸 더미 파일을 생성하는 유틸리티
async function createDummyLog() {
  const ws = createWriteStream(SOURCE_FILE);
  ws.write('192.168.0.1 - GET / HTTP/1.1 200\n');
  ws.write('10.0.0.2 - POST /login HTTP/1.1 500\n');
  ws.write('172.16.0.3 - GET /api/data HTTP/1.1 404\n');
  ws.write('192.168.0.4 - GET /assets/img.png HTTP/1.1 200\n');
  ws.write('10.0.0.5 - PUT /api/user HTTP/1.1 503\n');
  ws.end();
  
  // 쓰기 스트림이 완전히 닫히고 파일 저장이 끝날 때까지 대기
  await new Promise(resolve => ws.on('finish', resolve));
}

async function extractErrorLogs() {
  const readStream = createReadStream(SOURCE_FILE, { encoding: 'utf8' });
  const writeStream = createWriteStream(TARGET_FILE, { encoding: 'utf8' });
  
  // readline 인터페이스를 생성하여 바이트 청크를 줄 단위 텍스트로 분리
  const rl = createInterface({
    input: readStream,
    crlfDelay: Infinity
  });

  // 변환 스트림: readline에서 한 줄씩 넘겨받아 조건을 검사한다.
  const filterStream = new Transform({
    objectMode: true,
    transform(chunk, encoding, callback) {
      const line = chunk.toString();
      // 로그의 끝부분이 500번대 상태 코드인지 정규표현식으로 확인
      if (line.match(/ 5\d\d$/)) {
        // 조건에 맞으면 다음 쓰기 스트림으로 줄바꿈 기호를 덧붙여 전송
        this.push(line + '\n');
      }
      // 처리가 끝났음을 알림 (조건에 맞지 않는 줄은 그대로 버려짐)
      callback();
    }
  });

  // readline 객체는 비동기 반복자(Async Iterable)를 지원한다.
  // pipeline에 연결하기 위해 제너레이터 함수로 브릿지를 만든다.
  async function* lineGenerator() {
    for await (const line of rl) {
      yield line;
    }
  }

  console.log('로그 추출을 시작합니다.');
  
  try {
    // 흐름: 줄 생성기 -> 500번대 필터링 -> 파일 쓰기
    await pipeline(
      lineGenerator(),
      filterStream,
      writeStream
    );
    console.log('500번대 오류 로그 추출이 완료되었습니다.');
  } catch (error) {
    console.error('데이터 파이프라인 실행 중 오류 발생:', error.message);
  }
}

// 스크립트 실행 시작점
await createDummyLog();
await extractErrorLogs();

줄별 해설

  • crlfDelay: Infinity: 운영체제마다 다른 줄바꿈 문자(\r\n 또는 \n)를 일관되게 하나의 줄바꿈으로 인식하도록 돕는 설정이다. 파일이 어떤 환경에서 작성되었든 정확히 한 줄씩 잘라낸다.
  • objectMode: true: 변환 스트림 내부 옵션이다. 스트림은 기본적으로 바이트 데이터(버퍼)를 다루지만, objectMode를 켜면 버퍼 외에 자바스크립트 객체나 분리된 문자열을 원형 그대로 주고받을 수 있다.
  • this.push(...)와 callback(): 변환 스트림 안에서 데이터를 다음 단계로 보내고 싶다면 this.push()를 호출한다. 만약 this.push()를 부르지 않고 callback()만 실행하면 해당 데이터 청크는 걸러져 사라진다.
  • async function* lineGenerator(): node:stream의 파이프라인은 출발지로 비동기 반복 가능한(Iterable) 객체를 받을 수 있다. for await...of 구문을 사용해 readline이 생성한 문자열을 변환 스트림으로 차례차례 주입하는 역할을 한다.
  • await pipeline(...): 연결된 스트림 중 하나라도 권한 오류나 디스크 공간 부족 등으로 실패하면, pipeline은 열려 있던 파일 식별자(File Descriptor)를 안전하게 닫고 즉시 프라미스 거부(Reject)를 발생시킨다.

실행 결과

$ node extract-errors.mjs
로그 추출을 시작합니다.
500번대 오류 로그 추출이 완료되었습니다.

$ cat server-errors.log
10.0.0.2 - POST /login HTTP/1.1 500
10.0.0.5 - PUT /api/user HTTP/1.1 503

실무에서 자주 틀리는 것

과거의 pipe 메서드 사용과 오류 처리 누락

스트림 객체의 pipe() 메서드를 연결한 뒤 첫 번째 읽기 스트림에만 오류 핸들러를 등록하는 실수를 많이 범한다.

// 틀린 코드
const rs = createReadStream('data.txt');
const ws = createWriteStream('out.txt');

rs.pipe(ws);
rs.on('error', (err) => console.error('읽기 실패:', err));
// 쓰기 스트림에 대한 오류 처리가 빠져 있다.

이 상태에서 권한 문제로 쓰기 스트림(목적지)에서 오류가 발생하면, 잡히지 않은 예외(Unhandled Exception)로 처리되어 프로세스가 그 즉시 비정상 종료된다. pipeline 함수를 도입하면 코드가 간결해지고 예외 상황에서도 안전해진다.

// 고친 코드
import { pipeline } from 'node:stream/promises';

try {
  await pipeline(
    createReadStream('data.txt'),
    createWriteStream('out.txt')
  );
} catch (err) {
  console.error('스트림 오류 통합 처리:', err);
}

스트림 청크를 배열에 누적하는 안티 패턴

스트림을 도입하고도 data 이벤트에서 청크를 전역 배열에 차곡차곡 모아두었다가 end 이벤트에서 한 번에 조작하려는 시도가 종종 발생한다.

// 틀린 코드
const chunks = [];
readStream.on('data', (chunk) => {
  chunks.push(chunk); // 들어오는 청크를 계속 메모리에 쌓는다.
});

readStream.on('end', () => {
  const completeData = Buffer.concat(chunks).toString();
  // 결국 큰 파일 전체가 메모리에 적재되어 스트림의 이점이 사라진다.
});

이는 스트림을 사용하는 의미를 퇴색시킨다. 청크가 들어오는 즉시 소비하고 버리도록 흐름을 설계해야 한다. 데이터 가공이 필요하다면 변환 스트림(Transform)을 구현하여 중간에 끼워 넣는 것이 바람직하다.

// 고친 코드
const processStream = new Transform({
  transform(chunk, encoding, callback) {
    // 들어온 조각을 바로 변환하여 다음으로 넘긴다.
    this.push(chunk.toString().toUpperCase());
    callback();
  }
});

await pipeline(readStream, processStream, writeStream);

한눈에 보기

스트림의 특징 및 모듈 요약
개념 및 역할 설명 관련 내장 모듈 함수
읽기 스트림 큰 파일이나 수신 데이터를 지정된 바이트 크기의 조각으로 나누어 읽어들인다. fs.createReadStream
변환 스트림 데이터가 흐르는 길목에 위치하여 문자열 변경, 압축, 조건 필터링 등을 수행한다. new Transform(...)
안전한 파이프라인 데이터 파이프를 구축하고, 예기치 못한 중단 시 자원을 자동으로 닫아준다. stream/promises의 pipeline
줄 단위 텍스트 분리 바이트 스트림을 개행 기호 기준으로 잘라 비동기 반복자로 만들어 준다. node:readline

연습 문제

  1. 스트림을 사용하지 않고 fs.readFileSync로 5GB 크기의 파일을 읽으려고 시도할 때 일어나는 현상을 V8 엔진의 메모리 구조와 연관 지어 서술한다.
  2. 전통적인 pipe() 메서드 체이닝 대신 프라미스 기반의 pipeline() 유틸리티 함수를 사용하는 것이 권장되는 두 가지 주요 장점을 작성한다.
  3. node:readline 인터페이스가 반환하는 객체를 for await...of 구문과 함께 사용할 때 얻는 가독성과 동시성 제어 측면의 장점을 설명한다.

정답과 해설

  1. Node.js 구동의 핵심인 V8 자바스크립트 엔진은 하나의 프로세스에 기본적으로 할당 가능한 힙(Heap) 메모리의 한도가 2GB~4GB 내외로 제한되어 있다. 5GB 파일을 단일 버퍼나 문자열로 메모리에 올리려 시도하면 한도를 즉시 초과하여 ERR_STRING_TOO_LONG이나 메모리 할당 실패 오류를 출력하고 프로세스가 비정상 종료된다.
  2. 첫째, 여러 연결된 스트림 중 어느 한 곳에서라도 오류가 발생할 경우, 열려 있는 나머지 파일이나 네트워크 연결 스트림을 시스템이 자동으로 파괴(Destroy)하여 자원 유출을 막는다. 둘째, 프라미스 객체를 반환하므로 try/catch 블록 하나만으로 전체 데이터 전송 흐름에서 발생하는 예외를 간편하게 포착할 수 있다.
  3. 이벤트 리스너 콜백(on('line')) 방식 대신 비동기 반복문을 사용하면, 위에서 아래로 실행되는 동기적인 코드 작성 방식을 유지할 수 있어 가독성이 향상된다. 또한 특정 로그 줄을 읽어 데이터베이스에 기록해야 하는 상황일 때 반복문 내부에서 await db.save(line)을 호출함으로써, 디스크 입출력과 데이터베이스 전송의 속도를 자연스럽게 맞출 수 있다.

댓글 0

아직 댓글이 없습니다. 첫 댓글을 남겨 보세요.

댓글을 남기려면 로그인이 필요합니다.