initialSync Atlas Stream Processing 스트림 changeEvent 프로세서는 변경 스트림 테일링을 시작하기 전에 문서를 삽입한 것처럼 Atlas 컬렉션 의 기존 문서를 수집할 수 있습니다. 이제 단일 $source 단계에서 initialSync 한 번에 둘 이상의 컬렉션 에 대해 를 실행 수 있습니다.
이 페이지에서는 단일 컬렉션과 다중 컬렉션 동기화 특성을 비교하고, 다중 컬렉션 동기화 구성하는 방법을 안내하며, 이를 기반으로 구축된 클러스터 간 복제 파이프라인 보여줍니다.
단일 컬렉션 초기 동기화와의 비교
coll $source 단계의 필드 단일 컬렉션 이름 또는 컬렉션 이름 배열 허용합니다. 배열 에서 이름을 지정하는 모든 컬렉션 필드 에 지정한 데이터베이스 에 속해야 db 합니다. 단일 단계에서 둘 이상의 데이터베이스 의 컬렉션을 동기화할 수 없습니다. 제공하는 값에 따라 동기화 범위가 $source 결정됩니다.
범위 | $source 구성 | 행동 |
|---|---|---|
단일 컬렉션 |
| Atlas Stream Processing 명명된 컬렉션 동기화합니다. |
Explicit list |
| Atlas Stream Processing 목록의 각 컬렉션 나열한 순서대로 동기화합니다. |
다중 컬렉션 초기 동기화의 특성
다중 컬렉션 initialSync 작업에는 다음과 같은 특성이 있습니다.
목록 순서대로 컬렉션을 동기화합니다. 를
coll: ["a", "b", "c"]사용하면 Atlas Stream Processinga의 파티션을 보다 먼저, 의b파티션을b보다 먼저 배출하지만,c유출이 이후 컬렉션 의 파티션으로 유휴 용량 채우는 경우는 예외입니다.완료된 컬렉션을 다시 동기화하지 않고 재개합니다. Atlas Stream Processing 은 완료된 컬렉션을 추적하므로 다시 시작된 스트림 프로세서는 완료된 컬렉션을 건너뛰고 마지막으로 동기화된 문서 에서 진행 중인 모든 컬렉션 다시 시작합니다. 이미 내보낸 문서를 다시 전달할 수 있으므로 반복되는 문서를 멱등원으로 처리하다 하도록 싱크를 설계하세요.
전체 동기화 진행 상황을 노출합니다.
stats().stats.initialSync모든 대상 컬렉션의 동기화 진행 상황과 문서 수를 합산하여 보고합니다. 자세한 학습 은 초기 동기화 진행 상황 확인을 참조하세요.
고려해야 할 동작 차이
다중 컬렉션 initialSync은 스트림 프로세서를 구성하고 작동하는 방식에 영향을 미치는 방식에서 단일 컬렉션 동기화 와 다릅니다.
병렬처리는 컬렉션 별이 아닌 전역적으로 적용됩니다. 을(를) 설정하면
initialSync.parallelism컬렉션 별이 아닌 전체 대상 컬렉션 설정하다 에서 동시 파티션 읽기 수가 제한됩니다.캐치업 범위는 소스 범위를 따릅니다. 동기화 단계가 끝나면 Atlas Stream Processing 소스의 범위와 일치하는 변경 스트림 (전체 데이터베이스 소스의 경우 데이터베이스 수준 변경 스트림 , 전체 클러스터 소스의 경우 클러스터 수준 변경 스트림 , 전체 클러스터 소스의 경우 필터링된 변경 스트림 엽니다. 명시적 컬렉션 목록. 이를 별도로 구성할 필요가 없습니다.
컬렉션 동기화 실패는 전체 프로세서에 장애가 발생합니다. 단일 컬렉션 동기화 와 마찬가지로, Atlas Stream Processing 목록의 컬렉션 중 하나를 동기화하지 못하면 전체 스트림 프로세서가 실패합니다.
일부 소스 측 구성은 대상에 복제되지 않습니다.
initialSync인덱스, 고정 사이즈 컬렉션과 같은 컬렉션 옵션, time series 컬렉션, 유효성 검사기 또는 기본값 데이터 정렬 또는 뷰 정의를 복제하지 않습니다. 스트림 프로세서를 시작하기 전에 대상에 생성합니다. 또한 Atlas Stream Processing 스트림 프로세서가 시작된 후에 생성한 컬렉션을 검색하지 않습니다. 동기화할 컬렉션 설정하다 은 해당 점 에서 고정됩니다.
다중 컬렉션 초기 동기화 구성
다음 절차에서는 sample_analytics 데이터 세트에서customers 고객 프로필을 보관하는 컬렉션과 accounts 각 고객의 금융 계정을 보관하는 컬렉션을 사용합니다. 이 절차에서는 두 컬렉션을 한 Atlas cluster 에서 다른 클러스터로 복사한 다음 문서가 변경될 때마다 최신 상태로 유지하는 스트림 프로세서를 구성합니다.
전제 조건
이 절차를 완료하기 전에 다음이 필요합니다.
소스 클러스터와 대상 클러스터 모두에 Atlas 로 연결되어 있는 스트림 처리 작업 공간입니다. 연결을 추가하려면 연결 관리를 참조하세요.
나열하려는 컬렉션 수를 지원하는 프로세서 계층 입니다. 자세한 학습 은 리소스 할당을 참조하세요.
이 섹션의 명령을 클러스터에 대해 실행합니다. 연결하려면 mongosh 를 통해 클러스터에 연결을 참조하세요.
소스 컬렉션을 준비하려면 다음 단계를 완료하세요.
샘플 데이터 세트를 로드합니다.
이 절차에서는 customers accounts sample_analytics 데이터 세트의 및 컬렉션을 사용합니다. 샘플 데이터를 로드하는 방법을 학습 Atlas배포로 샘플 데이터 가져오기를 참조하세요.
변경 스트림 사전 및 게시 이미지를 활성화합니다.
절차의 단계에서 customers 를 설정하므로 및 accounts 컬렉션에서 사전 $source 및 사후 이미지를 fullDocument: "required" 활성화합니다. 클러스터 에 대해 다음 명령어를 실행합니다.
db.getSiblingDB("sample_analytics").runCommand({ collMod: "customers", changeStreamPreAndPostImages: { enabled: true } })
{ ok: 1, ... }
db.getSiblingDB("sample_analytics").runCommand({ collMod: "accounts", changeStreamPreAndPostImages: { enabled: true } })
{ ok: 1, ... }
에 대해 자세히 fullDocument: "required" 학습 MongoDB 컬렉션 변경 스트림을 참조하세요.
절차
Atlas UI 또는 를 선택하여 스트림 프로세서를 mongosh 구성합니다.
초기 동기화 진행 상황 확인
스트림 프로세서를 시작한 후 sp.<processor-name>.stats().stats.initialSync를 실행 모든 대상 컬렉션에서 전체 동기화 진행률을 확인합니다.
sp.replicate_analytics_sp.stats().stats.initialSync
{ progress: 1, numInputDocuments: Long('2246') }
progress 초기 동기화 완료되었는지 여부를 보고합니다. numInputDocuments은(는) Atlas Stream Processing 지금까지 대상 컬렉션에서 수집한 총 문서 수를 보고합니다.
예시
앞의 절차에서는 Atlas 클러스터 간에 두 개의 컬렉션을 복제합니다. 다음 예시 동일한 파이프라인 다른 컬렉션 설정하다 에 적용합니다. 단계의 db 및 coll 값만 $source 변경됩니다.
파이프라인 orders 및 customers의 기존 문서를 동기화한 다음 변경 사항이 발생하면 두 컬렉션에 변경 사항을 계속 복제합니다. 집계 세 단계로 이루어집니다.
$source단계는orderscustomers<source-connection-name>에서 와 를 동기화한 다음, 해당 두 컬렉션으로 범위가 지정된 변경 스트림 엽니다.$replaceRoot단계에서는 삭제 이벤트에 대해 문서 키를 선택하고 다른 모든 이벤트 에 대해 전체 문서 선택합니다. 메타데이터 문서 본문의 일부가 아니므로stream.source.*메타데이터 이 단계에서 유지됩니다.$merge단계는 의 일치하는 컬렉션 에 각 이벤트<destination-connection-name>작성하고, 소스 이벤트의 네임스페이스 메타데이터 사용하여 쓰기 (write) 를 라우팅하고 작업 유형을 선택하여 대상 문서 삽입, 교체 또는 삭제합니다.
sp = db.createStreamProcessor("replicate-app-db", [ { $source: { connectionName: "<source-connection-name>", db: "app", coll: ["orders", "customers"], config: { fullDocument: "required" }, initialSync: { enable: true } } }, { $replaceRoot: { newRoot: { $cond: { if: { $eq: [{ $meta: "stream.source.operationType" }, "delete"] }, then: { $meta: "stream.source.documentKey" }, else: "$fullDocument" } } } }, { $merge: { into: { connectionName: "<destination-connection-name>", db: { $meta: "stream.source.ns.db" }, coll: { $meta: "stream.source.ns.coll" } }, on: "_id", whenMatched: { $cond: { if: { $eq: [{ $meta: "stream.source.operationType" }, "delete"] }, then: "delete", else: "replace" } }, whenNotMatched: { $cond: { if: { $eq: [{ $meta: "stream.source.operationType" }, "delete"] }, then: "discard", else: "insert" } } } } ]); sp.start();
이 스트림 프로세서를 시작하기 전에 orders 및 customers에 필요한 인덱스, 컬렉션 옵션 또는 뷰를 <destination-connection-name>에서 생성합니다. initialSync는 이를 복제하지 않습니다.
추가 정보
이 가이드의 단계와 개념에 대해 자세히 알아보려면 다음 리소스를 참조하세요.
단계와 해당 필드에
$sourceinitialSync대해 학습 MongoDB 컬렉션 변경 스트림을 참조하세요.Atlas Stream Processing 진행 중인 동기화 를 어떻게 체크포인트하는지 학습 체크포인트를 참조하세요.
스트림 프로세서를 생성, 시작, 중지 및 모니터링하는 방법은 스트림 프로세서 개발 및 관리를 참조하세요.