
Atlas Stream Processing 및 Kinesis를 사용하여 MongoDB 변경 이벤트를 Amazon Redshift로 스트리밍하기
요약
MongoDB Atlas Stream Processing과 AWS Kinesis를 활용하여 MongoDB의 변경 이벤트를 Amazon Redshift로 실시간 스트리밍하는 방법을 설명합니다. S3나 Firehose 같은 중간 단계 없이 데이터 파이프라인을 단순화하여 운영 데이터를 분석 시스템으로 즉시 전달하는 아키텍처를 다룹니다.
핵심 포인트
- S3 스테이징이나 ETL 작업 없이 실시간 데이터 파이프라인 구축 가능
- Atlas Stream Processing을 통한 효율적인 CDC 및 데이터 재구성
- Kinesis를 활용한 데이터 내구성 확보 및 프로듀서-컨슈머 분리
- Redshift Streaming Ingestion을 통한 직접적인 데이터 쿼리 가능
이 튜토리얼은 Igor Alekseev에 의해 작성되었습니다.
운영 데이터를 분석 시스템으로 이동시키는 실시간 데이터 파이프라인 (Real-time data pipelines)은 현대 아키텍처의 초석입니다. 이 포스트에서는 중간에 S3 스테이징이나 Firehose를 거치지 않고, MongoDB Atlas Stream Processing과 AWS Kinesis Data Streams를 사용하여 MongoDB 변경 이벤트 (change events)를 Amazon Redshift로 스트리밍하는 방법을 살펴보겠습니다.
아키텍처 개요 (Architecture Overview)
흐름은 간단합니다:
- MongoDB Atlas — 컬렉션에서 변경 사항 (삽입, 업데이트, 교체)이 발생합니다.
- Atlas Stream Processing — 스트림 프로세서 (stream processor)가 변경 스트림 (change stream)을 통해 이러한 변경 사항을 캡처하고, 필터링 및 평탄한 JSON (flat JSON) 형태로 재구성하여 Kinesis로 레코드를 방출합니다.
- AWS Kinesis Data Stream — 소스와 분석 저장소 사이에서 내구성이 있고 순서가 보장되는 버퍼 (buffer) 역할을 합니다.
- Amazon Redshift — 스트리밍 구체화된 뷰 (streaming materialized view)가 Kinesis에서 직접 데이터를 읽어 쿼리 가능한 컬럼으로 데이터를 노출합니다.
여기서 핵심적인 장점은 단순함입니다. 중간 S3 버킷도, Firehose 전달 스트림 (delivery stream)도, 관리해야 할 ETL 작업도 없습니다. 데이터는 운영 쓰기 작업에서 분석 쿼리로 거의 실시간으로 흐릅니다.
왜 이 조합인가?
Atlas Stream Processing은 변경 데이터 캡처 (Change Data Capture, CDC)의 가장 어려운 부분을 처리합니다. MongoDB 변경 스트림에 네이티브하게 연결되고, 재개 토큰 (resume tokens)을 처리하며, 문서를 즉석에서 재구성하기 위한 파이프라인 언어 (집계 프레임워크, aggregation framework)를 제공합니다. 커스텀 컨슈머 (custom consumer)를 직접 작성하고 호스팅하는 대신, 파이프라인을 선언하기만 하면 Atlas가 실행을 관리합니다.
Kinesis는 프로듀서 (producer, Atlas)와 컨슈머 (consumer, Redshift)를 분리하는 내구성이 있고 순서가 보장된 스트림을 제공합니다. Redshift에 유지 관리 시간 (maintenance window)이 있거나 데이터를 다시 재생 (replay)해야 하는 경우, Kinesis는 레코드를 최대 365일 동안 보관합니다.
Redshift Streaming Ingestion은 데이터를 먼저 S3에 랜딩(landing)시키는 전통적인 패턴을 제거합니다. 구체화된 뷰 (Materialized View)가 Kinesis 스트림에서 직접 데이터를 읽어오므로, 새로운 레코드가 도착한 후 몇 초 이내에 쿼리가 가능해집니다.
구현 (Implementation)
1단계: Kinesis 데이터 스트림 생성
aws kinesis create-stream \
--stream-name atlas-change-events \
--shard-count 1 \
...
단일 샤드 (Shard)는 최대 1 MB/s 또는 초당 1,000개의 레코드 수집을 처리합니다. 변경 이벤트 볼륨에 따라 샤드를 확장하세요.
2단계: Atlas를 위한 IAM 설정
Atlas Stream Processing은 IAM AssumeRole을 통해 Kinesis에 인증합니다. Atlas AWS 계정이 역할을 맡을 수 있도록 허용하는 신뢰 정책 (Trust Policy)을 가진 역할을 생성하고, 최소한의 Kinesis 쓰기 권한을 부여하세요:
{
"Version": "2012-10-17",
"Statement": [
...
그게 전부입니다 — 단 두 가지 권한만 있으면 됩니다. Atlas는 배치 API (PutRecords)를 독점적으로 사용하므로, 단일 레코드 API (PutRecord)나 전체 DescribeStream 권한은 필요하지 않습니다.
3단계: Atlas에 Kinesis 연결 등록
Atlas UI에서 Stream Processing → Connection Registry → Add Connection → Kinesis로 이동하여 다음 정보를 제공합니다:
- 스트림 이름 (Stream name)
- AWS 리전 (AWS region)
- IAM 역할 ARN (IAM role ARN)
4단계: 스트림 프로세서 (Stream Processor) 생성
스트림 프로세서는 지속적으로 실행되는 집계 파이프라인 (Aggregation Pipeline)입니다. 패턴은 다음과 같습니다:
sp.createStreamProcessor("atlasToKinesisProcessor", [
// 1. 관련 작업 유형 필터링
{ $match: { operationType: { $in: ["insert", "update", "replace"] } } },
...
$emit 스테이지는 Kinesis 전용입니다. stream, region, partitionKey는 모두 최상위 필드이며, config 객체 내부에 중첩되지 않는다는 점에 유의하세요.
프로세서를 생성한 후, Stream Processing Instance (SPI)에 연결된 mongosh에서 프로세서를 시작하세요:
sp.atlasToKinesisProcessor.start();
5단계: Redshift Streaming Ingestion 설정
Redshift 측에서 Kinesis를 가리키는 외부 스키마 (External Schema)를 생성합니다:
CREATE EXTERNAL SCHEMA IF NOT EXISTS kinesis_schema
FROM KINESIS
IAM_ROLE 'arn:aws:iam::<YOUR_ACCOUNT>:role/<REDSHIFT_ROLE>';
그 다음, JSON 레코드를 파싱하는 스트리밍 구체화된 뷰 (Streaming Materialized View)를 생성합니다:
CREATE MATERIALIZED VIEW mdb_change_events AUTO REFRESH YES AS
SELECT
approximate_arrival_timestamp,
...
AUTO REFRESH YES 설정을 통해 Redshift는 새로운 레코드를 지속적으로 수집합니다. 필요한 경우 REFRESH MATERIALIZED VIEW mdb_change_events; 명령을 통해 수동으로 새로고침을 트리거할 수도 있습니다.
6단계: 데이터 쿼리하기
SELECT * FROM mdb_change_events
ORDER BY approximate_arrival_timestamp DESC
LIMIT 10;
이제 완료되었습니다. MongoDB의 운영 쓰기 (Operational writes) 데이터가 이제 몇 초 이내에 Redshift에서 쿼리 가능해집니다.
학습 내용 (Lessons Learned)
- 발행 전 평탄화 (Flatten before emitting). Redshift의
JSON_EXTRACT_PATH_TEXT는 중첩 수준이 5단계로 제한되어 있으며, 레코드 크기 제한이 16MB입니다. AWS는 대신 SUPER 타입을 사용하는JSON_PARSE로 전환할 것을 권장합니다. Atlas Stream Processing에서 ($replaceRoot+$project)를 사용하여 데이터를 평탄화하면, Kinesis는 단순한 평면 JSON (Flat JSON)을 수신하게 되고 Redshift의 구체화된 뷰 (Materialized view)를 깔끔하게 유지할 수 있습니다. - 구체화된 뷰는 새로운 레코드만 캡처합니다. 뷰가 생성될 당시 스트림에 이미 존재하던 레코드는 나타나지 않을 수 있습니다. 프로세서를 시작하기 전에 뷰를 생성하거나, 생성 후에 새로운 테스트 데이터를 삽입하십시오.
- 외부 스키마 생성에는 슈퍼유저 (Superuser) 권한이 필요합니다.
CREATE EXTERNAL SCHEMA문은 데이터베이스 소유자 또는 CREATE SCHEMA 권한을 가진 사용자가 실행해야 합니다. 만약 IAM 연동 세션에 이러한 권한이 없다면, 대신 관리자 사용자 및 비밀번호로 접속하십시오.
이 패턴을 사용해야 하는 경우
이 아키텍처는 다음과 같은 경우에 효과적입니다:
- 운영 중인 MongoDB 데이터에 대해 근실시간 (near-real-time) 분석이 필요한 경우
- CDC 인프라 (Debezium, Kafka Connect, 커스텀 컨슈머 (custom consumers))를 관리하는 것을 피하고 싶은 경우
- 분석 쿼리가 Redshift에서 실행되며, 컬럼형 (columnar) 성능을 원하는 경우
- 운영 오버헤드가 최소화된 서버리스 (serverless) 및 관리형 (managed) 방식을 선호하는 경우
더 높은 처리량 (throughput) 또는 더 복잡한 변환 (transformation)이 필요한 경우, 더 많은 Kinesis 샤드 (shards)를 추가하고 그에 따라 Atlas Stream Processing 인스턴스를 확장할 수 있습니다. 아키텍처는 동일하게 유지되며, 단지 용량 (capacity)만 변경됩니다.
참고 문헌 (References)
AI 자동 생성 콘텐츠
본 콘텐츠는 Dev.to AI tag의 원문을 AI가 자동으로 요약·번역·분석한 것입니다. 원 저작권은 원저작자에게 있으며, 정확한 내용은 반드시 원문을 확인해 주세요.
원문 바로가기