NestJS Interceptor가 Observable을 강제하는 AOP 설계 철학

miinhho·2025년 7월 30일
post-thumbnail

NestJS 는 Interceptor 에서 RxJS 의 Observable 을 반환하도록 해요. 이 부분은 NestJS 의 소스코드를 보면 더 드러나요.

// defer 연산자를 통해 Observable 반환
const nextFn = async (i = 0) => {
  if (i >= interceptors.length) {
    return defer(AsyncResource.bind(() => this.transformDeferred(next)));
  }
  const handler: CallHandler = {
    handle: () =>
    defer(AsyncResource.bind(() => nextFn(i + 1))).pipe(mergeAll()),
  };
  return interceptors[i].intercept(context, handler);
};
return defer(() => nextFn()).pipe(mergeAll());
}

이는 AOP(Aspect-Oriented Programming)의 본질적 특성을 웹 애플리케이션에서 완벽하게 구현하기 위한 아키텍처적 결정이에요.

Observable 이란?

Observable 은 기존의 관찰자 패턴을 시간축과 생명주기로 확장한 구조에요.

interface Observer<T> {
  next: (value: T) => void;     // 데이터 스트림 관찰
  error?: (error: any) => void; // 에러 상황 관찰  
  complete?: () => void;        // 완료 상황 관찰
}

// AOP에서 모든 실행 지점을 관찰 가능
const observable = new Observable(observer => {
  observer.next('시작');        // Join Point 1
  observer.next('진행 중');     // Join Point 2  
  observer.next('거의 완료');   // Join Point 3
  observer.complete();         // Join Point 4
});

Cross-cutting Concerns와 Join Point 제어

AOP가 해결하려는 문제

// 로깅
console.log('Creating user...');
// 인증 검증
if (!this.authService.isAuthorized()) throw new UnauthorizedException();
// 유효성 검사
if (!this.validateUserData(userData)) throw new BadRequestException();
// 성능 측정 시작
const startTime = Date.now();

try {
  // 성공 로깅
  console.log(`User created in ${Date.now() - startTime}ms`);
  // 이벤트 발행
  this.eventEmitter.emit('user.created', user);

  return user;
} catch (error) {
  // 에러 로깅
  console.error('User creation failed:', error);
  // 에러 변환
  throw new InternalServerErrorException();
}

AOP는 이런 횡단 관심사(Cross-cutting Concerns) 를 핵심 로직에서 분리하여, Join Point(실행 지점) 에서 투명하게 적용하는 것이 목표에요.

전통적 AOP

Spring AOP

@Around("execution(* com.example.service.*.*(..))")
public Object logExecutionTime(ProceedingJoinPoint joinPoint) throws Throwable {
    long startTime = System.currentTimeMillis();
    
    Object result = joinPoint.proceed(); // 동기적 실행
    
    long endTime = System.currentTimeMillis();
    logger.info("Execution time: " + (endTime - startTime) + "ms");
    
    return result;
}

동기적이고, 명확한 시작과 끝의 경계가 나타나요.

Promise 기반 AOP의 한계

@Injectable()
export class PromiseBasedInterceptor implements NestInterceptor {
  async intercept(context: ExecutionContext, next: CallHandler): Promise<any> {
    console.log('Before execution');
    const startTime = Date.now();
    
    try {
      const result = await next.handle().toPromise();
      
      // 1: 스트림의 중간값들을 놓침
      // SSE나 WebSocket에서 여러 번 전송되는 데이터를 어떻게 처리?
      
      console.log(`Execution time: ${Date.now() - startTime}ms`);
      return result;
      
    } catch (error) {
      // 2: 에러 복구나 재시도 로직을 어떻게 구현?
      // 3: 스트림이 중단된 경우와 에러의 구분은?
      console.error('Error occurred:', error);
      throw error;
    }
    
    // 4: 스트림이 완료된 후의 정리 작업은 언제?
    // 5: 클라이언트가 연결을 끊은 경우는?
  }
}

실제 문제 상황

1. SSE(Server-Sent Events)

@Get('events')
getEvents(): Observable<MessageEvent> {
  return interval(1000).pipe(
    map(i => ({ data: `Event ${i}` }))
  );
}

Promise 기반 인터셉터는 첫 번째 이벤트만 로깅하고 끝나요. 나머지 이벤트들은 추적이 불가능하죠.

2. 재시도가 필요한 API

@Get('external-data')
getExternalData(): Observable<ExternalData> {
  return this.httpService.get('/external-api').pipe(
    retry(3),
    timeout(5000),
    catchError(err => of({ error: 'Service unavailable' }))
  );
}

Promise 기반 인터셉터는 재시도 과정을 추적할 수 없어요.

Observable 이 제공하는 강력한 AOP

@Injectable()
export class ComprehensiveInterceptor implements NestInterceptor {
  intercept(context: ExecutionContext, next: CallHandler): Observable<any> {
    const requestId = this.generateRequestId();
    const startTime = Date.now();
    
    console.log(`[${requestId}] Request started`);
    
    return next.handle().pipe(
      // 실행 시작 시점 제어
      tap(() => console.log(`[${requestId}] First value emitted`)),
      
      // 각 데이터 포인트에서의 개입
      map(data => {
        console.log(`[${requestId}] Processing data:`, data);
        return this.transformData(data);
      }),
      
      // 조건부 로직 삽입
      filter(data => this.shouldEmitData(data)),
      
      // 에러 시점에서의 정교한 제어
      catchError(error => {
        console.error(`[${requestId}] Error occurred:`, error);
        
        // 에러 타입에 따른 다양한 처리
        if (error instanceof TimeoutError) {
          return this.handleTimeout(requestId);
        } else if (error instanceof HttpException) {
          return this.handleHttpError(error);
        }
        
        // 완전히 다른 Observable로 교체 가능
        return this.getFallbackData();
      }),
      
      // 재시도 로직
      retryWhen(errors => 
        errors.pipe(
          tap(err => console.log(`[${requestId}] Retrying after error:`, err)),
          delay(1000),
          take(3)
        )
      ),
      
      // 스트림 완료 시점 (성공/실패 무관)
      finalize(() => {
        const duration = Date.now() - startTime;
        console.log(`[${requestId}] Request completed in ${duration}ms`);
        this.cleanupResources(requestId);
      }),
      
      // 구독 해제 시점 (클라이언트 연결 끊김)
      tap({
        unsubscribe: () => {
          console.log(`[${requestId}] Client disconnected`);
          this.handleClientDisconnection(requestId);
        }
      })
    );
  }
}

AOP 에서 Observable의 핵심 가치

완벽한 Aspect Weaving

AOP 개념Observable 구현실제 효과
Join Point스트림의 각 연산자 지점실행 흐름의 모든 지점에 개입 가능
Pointcutpipe() 를 통한 선택적 적용조건부 로직 삽입
Advice연산자들(tap, map, catchError 등)Before/After/Around/AfterThrowing 모두 지원
WeavingRxJS 의 defer컴파일 타임이 아닌 런타임에 동적 조합

Aspect의 조합성 (Composability)

return next.handle().pipe(
  // 보안
  tap(() => this.securityService.logAccess()),
  // 성능 
  tap(() => this.metricsService.recordStart()),
  // 캐싱
  switchMap(data => this.cacheService.getCachedOrFetch(data)),
  // 변환
  map(data => this.transformationService.transform(data)),
  // 검증
  filter(data => this.validationService.validate(data)),
  // 로깅
  tap(data => this.loggingService.logSuccess(data)),
  // 에러 처리
  catchError(error => this.errorHandlingService.handle(error)),
  // 정리
  finalize(() => this.cleanupService.cleanup())
);

결론: 왜 Observable인가?

NestJS가 Observable 을 강제하는 이유는 AOP의 이상적인 구현을 위해서에요.

1. 시간축을 가진 관찰

  • 전통적인 관찰자 패턴: 단발성 이벤트만 관찰
  • Observable : 시간에 걸친 연속적 상태 변화 관찰

2. 투명한 횡단 관심사 적용

  • 비즈니스 로직은 Observable을 의식하지 않음
  • Interceptor가 투명하게 부가 기능을 주입

3. 조합 가능한 Aspect 설계

  • 각 연산자가 독립적인 Aspect
  • 선언적으로 Aspect를 조합하여 복잡한 로직 구성

4. 생명주기의 완전한 추상화

  • 시작, 진행, 완료, 에러, 취소를 모두 추상화
  • 모든 비동기 패턴을 동일하게 처리
@Injectable()
export class ModernAOPInterceptor implements NestInterceptor {
  intercept(context: ExecutionContext, next: CallHandler): Observable<any> {
    // - 비즈니스 로직은 전혀 수정하지 않으면서
    // - 모든 실행 시점에 횡단 관심사를 주입
    // - 다양한 비동기 패턴을 통일된 방식으로 처리
    // - 선언적이고 조합 가능한 방식으로 Aspect 구성
    
    return next.handle().pipe(
      // 여기서 AOP의 모든 패턴이 자연스럽게 구현되요
    );
  }
}
profile
재미있는 걸 좋아합니다

1개의 댓글

comment-user-thumbnail
2025년 7월 30일

깔끔하게 장단점이 잘 정리된 글이네요
잘보고갑니다. 🙏

답글 달기