deep-dive2025-02-28·9 min·164/348

MongoDB Change Streams 활용: 실시간 데이터 변경 감지

MongoDB Change Streams를 활용한 실시간 데이터 변경 감지 및 이벤트 기반 아키텍처 구현 방법을 다룹니다.

MongoDB Change Streams 활용

Introduction

MongoDB Change Streams는 컬렉션, 데이터베이스 또는 클러스터에서 발생하는 데이터 변경 사항을 실시간으로 감지할 수 있는 기능입니다. 전통적인 폴링 방식의 한계를 극복하고, 이벤트 기반 아키텍처를 구현할 수 있게 해줍니다. 이 글에서는 Change Streams의 내부 동작 원리와 실제 활용 사례를 살펴보겠습니다.

Environment

# MongoDB 버전 확인 (4.0 이상 필요)
mongosh --version

# 연결 테스트
mongosh "mongodb://localhost:27017"
// Node.js 환경
// package.json
{
  "name": "change-stream-demo",
  "dependencies": {
    "mongodb": "^6.3.0"
  }
}

Problem

기존 폴링 방식의 문제점:

// 기존 폴링 방식 - 비효율적
async function pollForChanges() {
  while (true) {
    const changes = await db.collection('orders').find({
      updatedAt: { $gte: lastChecked }
    }).toArray();
    
    if (changes.length > 0) {
      processChanges(changes);
      lastChecked = new Date();
    }
    
    // 1초마다 폴링 - 리소스 낭비
    await new Promise(resolve => setTimeout(resolve, 1000));
  }
}

이 방식의 문제점:

  • 데이터베이스에 불필요한 부하 발생
  • 변경 감지 지연 시간(latency) 불가피
  • 네트워크 대역폭 낭비

Analysis

Change Streams의 동작 원原理를 분석했습니다:

// Change Streams 내부 동작
// oplog 기반으로 작동
const pipeline = [
  {
    $match: {
      "operationType": { $in: ["insert", "update", "replace"] },
      "fullDocument.status": "pending"
    }
  },
  {
    $project: {
      _id: 0,
      operationType: 1,
      documentKey: 1,
      fullDocument: 1,
      updateDescription: 1,
      ns: 1,
      clusterTime: 1
    }
  }
];

// Resume token 관리
let resumeToken = null;

const changeStream = db.collection('orders').watch(pipeline, {
  resumeAfter: resumeToken,
  fullDocument: 'updateLookup'
});

Solution

기본 Change Streams 구현

// change-stream.js
const { MongoClient } = require('mongodb');

const uri = 'mongodb://localhost:27017';
const client = new MongoClient(uri);

async function watchChanges() {
  try {
    await client.connect();
    const database = client.db('myapp');
    const collection = database.collection('orders');

    // Change Streams 시작
    const changeStream = collection.watch([], {
      fullDocument: 'updateLookup'
    });

    console.log('Listening for changes...');

    changeStream.on('change', (change) => {
      console.log('Change detected:', {
        operationType: change.operationType,
        documentKey: change.documentKey,
        fullDocument: change.fullDocument,
        timestamp: change.clusterTime
      });

      // 이벤트 처리 로직
      switch (change.operationType) {
        case 'insert':
          handleInsert(change.fullDocument);
          break;
        case 'update':
          handleUpdate(change.fullDocument, change.updateDescription);
          break;
        case 'delete':
          handleDelete(change.documentKey._id);
          break;
      }
    });

    // 에러 처리
    changeStream.on('error', (error) => {
      console.error('Change stream error:', error);
      // 재연결 로직
      setTimeout(watchChanges, 5000);
    });

  } catch (error) {
    console.error('Connection error:', error);
  }
}

watchChanges();

Resume Token 관리

// resume-token-manager.js
const fs = require('fs');

class ResumeTokenManager {
  constructor(tokenPath = './resumeToken.json') {
    this.tokenPath = tokenPath;
    this.token = this.loadToken();
  }

  loadToken() {
    try {
      if (fs.existsSync(this.tokenPath)) {
        const data = fs.readFileSync(this.tokenPath, 'utf8');
        return JSON.parse(data).resumeToken;
      }
    } catch (error) {
      console.error('Error loading token:', error);
    }
    return null;
  }

  saveToken(token) {
    this.token = token;
    fs.writeFileSync(this.tokenPath, JSON.stringify({
      resumeToken: token,
      updatedAt: new Date().toISOString()
    }));
  }

  getResumeToken() {
    return this.token;
  }
}

// 사용 예
const tokenManager = new ResumeTokenManager();

const changeStream = collection.watch([], {
  resumeAfter: tokenManager.getResumeToken()
});

changeStream.on('change', (change) => {
  // 변경 처리
  processChange(change);
  
  // 토큰 저장
  tokenManager.saveToken(change._id);
});

필터링된 Change Streams

// 특정 조건의 변경만 감지
const filteredPipeline = [
  {
    $match: {
      "fullDocument.priority": "high",
      "operationType": { $in: ["insert", "update"] }
    }
  }
];

const changeStream = collection.watch(filteredPipeline, {
  fullDocument: 'updateLookup'
});

// 복수 컬렉션 감지
const database = client.db('myapp');
const changeStream = database.watch([], {
  fullDocument: 'updateLookup'
});

Lessons Learned

  1. Replica Set 필수: Change Streams는 Replica Set이나 Sharded Cluster에서만 사용 가능합니다
  2. Resume Token 관리: 반드시 이전 토큰을 저장하고 복구 시 활용해야 합니다
  3. 필터링 활용: $match 파이프라인으로 불필요한 이벤트를 필터링하면 성능이 향상됩니다
  4. 에러 핸들링 필수: 연결 끊김에 대한 재연결 로직을 반드시 구현해야 합니다
  5. 드라마틱 이벤트: invalidate 이벤트로 컬렉션 이름 변경이나 드롭을 감지할 수 있습니다

이 블로그는 외부 스폰서십, 제휴 마케팅 또는 광고 수익을 받지 않습니다.