troubleshooting2025-02-22·8 min·178/348

MongoDB Change Stream 연결 끊김: 재연결 및 복구 전략

MongoDB Change Stream 연결 끊김 문제를 해결하는 방법과 안정적인 재연결 메커니즘을 구현하는 방법을 설명합니다.

MongoDB Change Stream 연결 끊김

Introduction

MongoDB Change Stream은 실시간 데이터 변경 감지에 유용하지만, 네트워크 불안정성이나 서버 문제로 인해 연결이 끊길 수 있습니다. 연결 끊김 시 이벤트를 놓치지 않도록 안정적인 재연결 메커니즘을 구현하는 것이 중요합니다. 이 글에서는 Change Stream 연결 끊김 문제와 해결 방법을 다루겠습니다.

Environment

// MongoDB 연결 설정
const { MongoClient } = require('mongodb');

const uri = 'mongodb://localhost:27017';
const client = new MongoClient(uri, {
  autoReconnect: true,
  reconnectTries: Number.MAX_VALUE,
  reconnectInterval: 1000
});
// 연결 상태 확인
mongosh --eval "db.serverStatus().connections"

Problem

Change Stream 연결 끊김 상황:

// 기본 Change Stream 구현
const collection = client.db('mydb').collection('events');
const changeStream = collection.watch();

changeStream.on('change', (change) => {
  console.log('Change:', change);
});

changeStream.on('error', (error) => {
  console.error('Change stream error:', error);
  // 에러 발생 시 연결 끊김
});

// 출력:
// Change stream error: MongoNetworkError: connection 5 to localhost:27017 closed
// Change stream error: MongoServerClosedError: Server closed the connection
# 네트워크 문제 확인
ping localhost
# PING localhost (127.0.0.1): 56 data bytes
# Request timeout for icmp_seq 0

# MongoDB 프로세스 상태 확인
ps aux | grep mongod
# mongod may not be running

Analysis

연결 끊김 원인을 분석했습니다:

// 연결 상태 모니터링
client.on('connectionPoolCreated', (event) => {
  console.log('Connection pool created:', event);
});

client.on('connectionCheckedOut', (event) => {
  console.log('Connection checked out:', event);
});

client.on('connectionCheckedIn', (event) => {
  console.log('Connection checked in:', event);
});

client.on('connectionPoolClosed', (event) => {
  console.log('Connection pool closed:', event);
});

// Change Stream 상태 확인
const changeStream = collection.watch();
console.log('Change stream is:', changeStream.isClosed());
// Resume Token 확인
let resumeToken = null;

changeStream.on('change', (change) => {
  resumeToken = change._id;
  console.log('Resume token updated:', resumeToken);
});

// 연결 끊김 후 재시작 시 Resume Token 필요

Solution

1단계: 안정적인 Change Stream 구현

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

class ResilientChangeStream {
  constructor(uri, dbName, collectionName, options = {}) {
    this.uri = uri;
    this.dbName = dbName;
    this.collectionName = collectionName;
    this.resumeToken = null;
    this.isRunning = false;
    this.retryDelay = options.retryDelay || 5000;
    this.maxRetries = options.maxRetries || 10;
    this.retryCount = 0;
  }

  async start() {
    this.isRunning = true;
    await this.connectAndWatch();
  }

  async connectAndWatch() {
    while (this.isRunning && this.retryCount < this.maxRetries) {
      let client;
      try {
        client = new MongoClient(this.uri, {
          autoReconnect: true,
          reconnectTries: Number.MAX_VALUE,
          reconnectInterval: 1000
        });

        await client.connect();
        console.log('Connected to MongoDB');

        const collection = client.db(this.dbName).collection(this.collectionName);
        
        const options = {};
        if (this.resumeToken) {
          options.resumeAfter = this.resumeToken;
        }

        const changeStream = collection.watch([], options);
        
        changeStream.on('change', (change) => {
          this.retryCount = 0;
          this.resumeToken = change._id;
          this.handleChange(change);
        });

        changeStream.on('error', async (error) => {
          console.error('Change stream error:', error);
          await this.close(changeStream, client);
          await this.waitForRetry();
        });

        // 연결 유지
        await new Promise((resolve) => {
          changeStream.on('close', resolve);
        });

      } catch (error) {
        console.error('Connection error:', error);
        if (client) {
          await client.close().catch(() => {});
        }
        await this.waitForRetry();
      }
    }
  }

  async waitForRetry() {
    this.retryCount++;
    console.log(`Retrying in ${this.retryDelay}ms (attempt ${this.retryCount}/${this.maxRetries})`);
    await new Promise(resolve => setTimeout(resolve, this.retryDelay));
  }

  async close(changeStream, client) {
    try {
      await changeStream.close();
    } catch (error) {
      // 무시
    }
    try {
      await client.close();
    } catch (error) {
      // 무시
    }
  }

  handleChange(change) {
    console.log('Change detected:', {
      operationType: change.operationType,
      documentKey: change.documentKey,
      timestamp: new Date()
    });
    
    // 비즈니스 로직 처리
    this.processChange(change);
  }

  async processChange(change) {
    // 구체적인 처리 로직
    switch (change.operationType) {
      case 'insert':
        await this.handleInsert(change.fullDocument);
        break;
      case 'update':
        await this.handleUpdate(change.fullDocument, change.updateDescription);
        break;
      case 'delete':
        await this.handleDelete(change.documentKey._id);
        break;
    }
  }

  async handleInsert(document) {
    console.log('New document:', document);
  }

  async handleUpdate(document, updateDescription) {
    console.log('Updated document:', document, 'Changes:', updateDescription);
  }

  async handleDelete(documentId) {
    console.log('Deleted document:', documentId);
  }

  stop() {
    this.isRunning = false;
  }
}

// 사용 예
const changeStream = new ResilientChangeStream(
  'mongodb://localhost:27017',
  'mydb',
  'events',
  {
    retryDelay: 5000,
    maxRetries: 20
  }
);

changeStream.start().catch(console.error);

// Graceful shutdown
process.on('SIGTERM', () => {
  changeStream.stop();
});

2단계: Resume Token 관리

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

class ResumeTokenStore {
  constructor(filePath = './resume-tokens.json') {
    this.filePath = filePath;
    this.tokens = this.loadTokens();
  }

  loadTokens() {
    try {
      if (fs.existsSync(this.filePath)) {
        const data = fs.readFileSync(this.filePath, 'utf8');
        return JSON.parse(data);
      }
    } catch (error) {
      console.error('Error loading tokens:', error);
    }
    return {};
  }

  saveTokens() {
    try {
      fs.writeFileSync(this.filePath, JSON.stringify(this.tokens, null, 2));
    } catch (error) {
      console.error('Error saving tokens:', error);
    }
  }

  getToken(collectionName) {
    return this.tokens[collectionName] || null;
  }

  setToken(collectionName, token) {
    this.tokens[collectionName] = token;
    this.saveTokens();
  }
}

// 사용 예
const tokenStore = new ResumeTokenStore();

const changeStream = collection.watch([], {
  resumeAfter: tokenStore.getToken('events')
});

changeStream.on('change', (change) => {
  tokenStore.setToken('events', change._id);
  console.log('Token saved:', change._id);
});

Lessons Learned

  1. Resume Token 필수: Change Stream은 Resume Token을 사용하여 중단된 위치부터 다시 시작할 수 있습니다
  2. 자동 재연결 구현: 네트워크 불안정성에 대비한 자동 재연결 로직이 반드시 필요합니다
  3. 재시도 제한: 무한 재시도보다는 지수 백오프와 함께 최대 재시도 횟수를 설정하는 것이 안전합니다
  4. Graceful Shutdown: 프로세스 종료 시 Change Stream을 정상적으로 닫아야 합니다
  5. 토큰 영속성: Resume Token은 파일이나 데이터베이스에 저장하여 프로세스 재시작 시에도 유지되어야 합니다

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