스트리밍 구체화된 뷰는 Atlas Stream Processing 스트림 프로세서가 현재 상태를 유지하는 컬렉션입니다. 이를 위해 스트림 프로세서는 소스 컬렉션의 각 변경 사항을 읽고, 변경의 영향을 계산하고, 결과를 뷰에 적용합니다. 이 뷰는 수동 또는 예정된 새로 고침 없이 소스 데이터의 현재 상태를 반영합니다.
스트리밍 구체화된 뷰는 데이터 소스가 변경되는 경우보다 계산된 결과를 읽는 워크로드(예: 대시보드, 누계 합계, 미결 항목 개수) 에 적합합니다. 데이터가 도착하면 뷰가 업데이트되므로 읽기는 지연 시간이 즐어들어 현재 결과를 반환합니다.
온디맨드 구체화된 뷰 비교
MongoDB는 구체화된 뷰에 대한 두 가지 접근 방식을 지원합니다.
온디맨드 구체화된 뷰 는 수동으로 또는 예정에 따라 실행하는 집계 파이프라인의 결과를 저장합니다.
스트리밍 구체화된 뷰는 연속적으로 실행되는 Atlas Stream Processing 스트리밍 프로세서의 결과를 저장합니다.
두 개 방법 모두 계산된 결과를 디스크에 저장하고 보기에서 직접 읽기를 제공합니다. 보기가 업데이트되는 방법과 시기가 다릅니다.
특성 | 온디맨드 구체화된 뷰 | 스트리밍 구체화된 뷰 |
|---|---|---|
trigger 업데이트 | 수동 또는 예정 | 변경 중심의 연속적 |
지연 시간 | 분을 일수로 변환 | 초 미만에서 수 초 단위 |
데이터 신선도 | 특정 점 스냅샷 | 영구적으로 동기화됨 |
컴퓨트 모델 | 전체 결과 다시 계산 | 증분 효과 계산 |
Best fit | 배치 보고, 주기적 집계 | 실시간 대시보드, 운영 분석 |
스트리밍 구체화된 뷰 프로세서의 특징
스트리밍 구체화된 뷰를 유지하는 것은 표준 Atlas Stream Processing 집계 단계로 구성하는 패턴입니다. 이 패턴을 구현하는 스트리밍 프로세서에는 다음과 같은 특성이 있을 수 있습니다.
변경 스트림 소스를 읽습니다.
$source단계는fullDocument및fullDocumentBeforeChange가required(으)로 설정된 소스 컬렉션에서 읽어오므로 파이프라인은 변경 전후 각 문서의 상태를 비교할 수 있습니다.이벤트별 부호 델타를 계산합니다.
$addFields단계는$switch표현식을 사용하여 변경이 계산된 결과에 영향을 미치는 방식에 따라 각 삽입, 업데이트 또는 삭제에 양수 또는 음수 값을 할당할 수 있습니다.창에서 결과를 그룹화합니다. 스트림 프로세서는 무제한 스트림에서 작동하므로 모든
$group단계는 창 단계 내에서 실행되어야 합니다. 창$group는 창 간격 동안 각 키의 델타를 합산합니다. 창 간격은 보기의 새로 고침 간격으로도 작용하므로 보기의 새로 고침 정도를 결정합니다. 설정할 수 있는 가장 작은 간격은 1밀리초입니다.결과를 가산적으로 적용합니다.
$merge단계는 각 창의 결과를 대체하는 대신 뷰의 실행 중 총계에 추가하는whenMatched파이프라인을 사용할 수 있습니다.싱크 단계에서 종료됩니다. 스트리밍 프로세서 파이프라인은 싱크 단계에서 종료될 수 있습니다. Atlas 컬렉션에 쓰기 (write) 위해
$merge를 사용합니다.
증분 집계만 스트리밍 구체화된 뷰로 변환됩니다. 전체 컬렉션 스캔이 필요한 집계는 자격이 없습니다.
고려해야 할 동작 차이
스트리밍 구체화된 뷰는 하위 소비자에 영향을 미치는 방식으로 배치 집계와 다르게 동작합니다.
Atlas UI에서 스트리밍 구체화된 뷰 만들기
다음 절차는 구매 방식에 따라 완료된 판매 건수를 유지하는 스트리밍 구체화된 뷰를 만듭니다. 스트림 프로세서는 sample_supplies.sales 컬렉션의 변경 스트림을 읽고 sample_supplies.sales_by_channel 컬렉션에 쓰기 (write)를 수행합니다.
참고
Stream Processing은 Apache Kafka 주제에서 읽고 Apache Iceberg 테이블에 AWS S3에 쓰기 (write)도 할 수 있습니다. 자세한 내용은 Apache Kafka Broker 및 $iceberg 집계 단계를 참조하세요.
전제 조건
Stream Processor를 만들기 전에 소스 데이터가 있는 클러스터에 대한 Atlas 연결이 있는 Stream Processing 작업 공간 이 있어야 합니다. 연결을 추가하려면 연결 관리를 참조하세요.
다음 단계를 완료하여 소스 컬렉션을 준비하고 보기를 시딩합니다.
샘플 데이터 세트를 로드합니다.
이 절차에서는 sample_supplies 데이터 세트의 sales 컬렉션을 사용합니다. 샘플 데이터를 로드하는 방법은 Atlas 배포서버에 샘플 데이터 가져오기를 참조하세요.
현재 결과를 사용하여 보기를 시드합니다.
클러스터에 대해 다음 배치 집계를 실행하여 현재 카운트로 뷰를 채우세요.
db.sales.aggregate([ { $match: { status: "completed" } }, { $group: { _id: "$purchaseMethod", active_count: { $sum: 1 } } }, { $merge: { into: "sales_by_channel", whenMatched: "replace", whenNotMatched: "insert" } } ])
시딧된 카운트를 확인하려면 보기를 쿼리합니다. sales_by_channel 컬렉션에는 구매 방식당 하나의 문서가 포함되어 있습니다.
db.sales_by_channel.find()
[ { _id: 'Phone', active_count: 596 }, { _id: 'Online', active_count: 1585 }, { _id: 'In store', active_count: 2819 } ]
이 시드 집계는 그 자체가 온디맨드 구체화된 뷰입니다. 두 가지 뷰 유형은 상호 보완적입니다. 배치 집계로 온디맨드 구체화된 뷰를 초기화한 다음 동일한 컬렉션을 현재 상태로 유지하는 스트림 프로세서를 시작할 수 있습니다. 온디맨드 뷰가 스트리밍 구체화된 뷰가 됩니다.
참고
파이프라인에 창 단계가 사용되지 않는 경우 대신 initialSync 으로 뷰를 시딩할 수 있습니다. 이 경우 스트림 프로세서는 먼저 소스 컬렉션에 있는 모든 문서를 삽입 이벤트로 수집한 다음 새로운 변경 이벤트를 처리합니다. $source 옵션 initialSync에 대해 자세히 알아보려면 MongoDB 컬렉션 변경 스트림을 참조하십시오.
절차
소스 구성.
Source 필드에서 Connection 드롭다운 목록에서 소스 클러스터에 대한 Atlas 연결을 선택합니다.
JSON 텍스트 상자에서
$source단계를 구성하여 사전 및 게시 이미지로sales컬렉션을 읽습니다.
{ "$source": { "connectionName": "<connection-name>", "db": "sample_supplies", "coll": "sales", "config": { "fullDocument": "required", "fullDocumentBeforeChange": "required" } } }
각 이벤트에 대한 변경 델타를 계산하는 단계를 추가합니다.
Start building your pipeline 창에서 + Custom stage을 클릭합니다.
JSON 텍스트 상자에서 각 삽입, 업데이트 또는 삭제에 서명된 델타를 할당하고 구매 방법을 캡처하는
$addFields단계를 추가합니다.
{ "$addFields": { "_delta": { "$switch": { "branches": [ { "case": { "$and": [ { "$eq": ["$operationType", "insert"] }, { "$eq": ["$fullDocument.status", "completed"] } ] }, "then": 1 }, { "case": { "$and": [ { "$eq": ["$operationType", "update"] }, { "$eq": ["$fullDocumentBeforeChange.status", "completed"] }, { "$eq": ["$fullDocument.status", "returned"] } ] }, "then": -1 }, { "case": { "$and": [ { "$eq": ["$operationType", "delete"] }, { "$eq": ["$fullDocumentBeforeChange.status", "completed"] } ] }, "then": -1 } ], "default": 0 } }, "_channel": { "$ifNull": [ "$fullDocument.purchaseMethod", "$fullDocumentBeforeChange.purchaseMethod" ] } } }
효과가 없는 이벤트를 제거하는 단계를 추가합니다.
+을 클릭한 다음 Custom stage을 선택합니다.
JSON 텍스트 상자에서 델타가 0인 이벤트를 제거하는
$match단계를 추가합니다.
{ "$match": { "_delta": { "$ne": 0 } } }
창 그룹 단계를 추가합니다.
+을 클릭한 다음 Custom stage을 선택합니다.
JSON 텍스트 상자에서 구매 방법별 델타를 1초 간격으로 합산하는
$tumblingWindow단계를 추가합니다.
{ "$tumblingWindow": { "boundary": "processingTime", "interval": { "size": 1, "unit": "second" }, "pipeline": [ { "$group": { "_id": "$_channel", "active_count": { "$sum": "$_delta" }, "windowStart": { "$first": { "$meta": "stream.window.start" } } } } ] } }
중요
스트리밍 프로세서에는 모든 $group 단계가 창 단계 내에서 실행되어야 합니다.
싱크를 구성합니다.
Sink 필드에서 Connection 드롭다운 목록에서 사용자의 Atlas 연결을 선택합니다.
JSON 텍스트 상자에서
$merge단계를 구성하여 각 창의 결과를sales_by_channel의 누적 합계에 추가합니다.
{ "$merge": { "into": { "connectionName": "<connection-name>", "db": "sample_supplies", "coll": "sales_by_channel" }, "whenMatched": [ { "$set": { "active_count": { "$cond": [ { "$gt": ["$$new.windowStart", { "$ifNull": ["$lastWindowStart", { "$toDate": 0 }] }] }, { "$add": ["$active_count", "$$new.active_count"] }, "$active_count" ] }, "lastWindowStart": { "$max": [ { "$ifNull": ["$lastWindowStart", { "$toDate": 0 }] }, "$$new.windowStart" ] } } } ], "whenNotMatched": "insert" } }
참고
lastWindowStart 최고 수위는 재생된 창이 이중 계산되는 것을 방지합니다.
프로세서 세부 정보를 입력합니다.
Stream processor name 필드에
sales_stats_sp를 입력합니다.Stream Processing의 계층을 선택합니다. 워크로드에 맞는 계층을 선택하려면 Atlas Stream Processing 계층 선택 가이드를 참조하세요.
스트림 프로세서를 시작합니다.
Stream Processors 탭에서 sales_stats_sp 을 선택하고 Start를 클릭합니다.
프로세서는 이제 sales_by_channel 를 계속 유지합니다. 스트림 프로세서의 시작, 중지 및 모니터링에 대한 자세한 내용은 스트림 프로세서 개발을 참조하십시오.
보기가 현재 상태를 유지하는지 확인합니다.
프로세서를 시작하면 sales 컬렉션의 변경 사항이 수초 내에 sales_by_channel 에 업데이트됩니다. 이를 확인하려면 새로운 온라인 판매를 완료하여 삽입합니다.
db.sales.insertOne({ saleDate: new Date(), purchaseMethod: "Online", status: "completed", items: [], customer: {}, couponUsed: false })
프로세서는 Online 카운트를 증가시킵 그리고 lastWindowStart에 창 경계를 기록합니다:
db.sales_by_channel.find()
[ { _id: 'Phone', active_count: 596 }, { _id: 'Online', active_count: 1586, lastWindowStart: ISODate('2026-07-23T15:18:01.000Z') }, { _id: 'In store', active_count: 2819 } ]
고객이 해당 판매를 반환하면 Online 건수가 시드된 값으로 돌아가답니다.
db.sales.updateOne( { purchaseMethod: "Online", status: "completed" }, { $set: { status: "returned" } } )
db.sales_by_channel.find()
[ { _id: 'Phone', active_count: 596 }, { _id: 'Online', active_count: 1585, lastWindowStart: ISODate('2026-07-23T15:18:13.000Z') }, { _id: 'In store', active_count: 2819 } ]
예시
이 예시들은 Atlas Stream Processing 예시 리포지토리의 smv 예시 를 살펴보는 과정입니다. 우선 순위 수준별 문서 하나씩으로 미결 지원 티켓의 누계 개수를 저장하는 queue_stats 컬렉션을 유지합니다. 스트리밍 프로세서는 support_tickets 변경 스트림을 읽고 티켓이 열리고, 해결되고, 삭제될 때 각 카운트를 업데이트합니다.
파이프라인은 티켓이 열릴 때 +1 로, 티켓이 해결되거나 삭제될 때 -1 로 티켓의 우선 순위에 따라 카운트를 조정합니다. 집계에는 다섯 단계가 있습니다.
$source단계는 사전 이미지와 게시 이미지로support_tickets변경 스트림을 읽습니다.$addFields단계는operationType에서$switch를 사용하여_delta을 계산하고_priority을 그룹 키로 추출합니다.$match단계는 델타가 0인 이벤트를 제거합니다.$tumblingWindow단계는 각 1초 창 내에서 우선 순위별 델타를 합산합니다.$merge단계는lastWindowStart를 고수위로 사용하여 재생 시 이중 계산을 피하기 위해 각 창의 델타를 실행 중인 총계에 추가합니다.
[ { "$source": { "connectionName": "<connection-name>", "db": "support", "coll": "support_tickets", "config": { "fullDocument": "required", "fullDocumentBeforeChange": "required" } } }, { "$addFields": { "_delta": { "$switch": { "branches": [ { "case": { "$and": [ { "$eq": ["$operationType", "insert"] }, { "$eq": ["$fullDocument.status", "open"] } ] }, "then": 1 }, { "case": { "$and": [ { "$eq": ["$operationType", "update"] }, { "$eq": ["$fullDocumentBeforeChange.status", "open"] }, { "$eq": ["$fullDocument.status", "resolved"] } ] }, "then": -1 }, { "case": { "$and": [ { "$eq": ["$operationType", "delete"] }, { "$eq": ["$fullDocumentBeforeChange.status", "open"] } ] }, "then": -1 } ], "default": 0 } }, "_priority": { "$ifNull": [ "$fullDocument.priority", "$fullDocumentBeforeChange.priority" ] } } }, { "$match": { "_delta": { "$ne": 0 } } }, { "$tumblingWindow": { "boundary": "processingTime", "interval": { "size": 1, "unit": "second" }, "pipeline": [ { "$group": { "_id": "$_priority", "open_count": { "$sum": "$_delta" }, "windowStart": { "$first": { "$meta": "stream.window.start" } } } } ] } }, { "$merge": { "into": { "connectionName": "<connection-name>", "db": "support", "coll": "queue_stats" }, "whenMatched": [ { "$set": { "open_count": { "$cond": [ { "$gt": ["$$new.windowStart", { "$ifNull": ["$lastWindowStart", { "$toDate": 0 }] }] }, { "$add": ["$open_count", "$$new.open_count"] }, "$open_count" ] }, "lastWindowStart": { "$max": [ { "$ifNull": ["$lastWindowStart", { "$toDate": 0 }] }, "$$new.windowStart" ] } } } ], "whenNotMatched": "insert" } } ]
queue_stats 컬렉션의 각 문서는 다음과 같습니다.
{ _id: "P1", open_count: <num>, lastWindowStart: <timestamp> }
하나의 이벤트가 두 개의 그룹 키를 조정해야 하는 경우 스캘라 _delta 를 _adjustments 배열로 대체하고 배열을 조정당 하나의 문서로 팬아웃합니다. 단일 키 파이프라인을 다음과 같이 수정합니다.
영향을 받은
$addFields그룹$switch_adjustments키 하나당 하나의 요소가 있는 배열을 각 분기가 반환하도록 단계를 대체합니다. 에스칼레이션은 두 개의 요소를 반환합니다.
단일 키 파이프라인의 2 및 3 단계를 다음과 같이 바꾸십시오.
// Stage 2 (replacement): Compute an _adjustments array. { $addFields: { _adjustments: { $switch: { branches: [ { case: { $and: [ { $eq: ["$operationType", "insert"] }, { $eq: ["$fullDocument.status", "open"] } ]}, then: [{ _priority: "$fullDocument.priority", _delta: 1 }] }, { case: { $and: [ { $eq: ["$operationType", "update"] }, { $eq: ["$fullDocumentBeforeChange.status", "open"] }, { $eq: ["$fullDocument.status", "resolved"] } ]}, then: [{ _priority: "$fullDocumentBeforeChange.priority", _delta: -1 }] }, { case: { $and: [ { $eq: ["$operationType", "update"] }, { $eq: ["$fullDocument.status", "open"] }, { $ne: ["$fullDocument.priority", "$fullDocumentBeforeChange.priority"] } ]}, then: [ { _priority: "$fullDocumentBeforeChange.priority", _delta: -1 }, { _priority: "$fullDocument.priority", _delta: 1 } ] }, { case: { $and: [ { $eq: ["$operationType", "delete"] }, { $eq: ["$fullDocumentBeforeChange.status", "open"] } ]}, then: [{ _priority: "$fullDocumentBeforeChange.priority", _delta: -1 }] } ], default: [] } } } }, // Stage 3 (replacement): Fan out into one document per adjustment, // then lift the adjustment fields back to the top level. { $unwind: "$_adjustments" }, { $set: { _priority: "$_adjustments._priority", _delta: "$_adjustments._delta" } }
추가 정보
이 가이드의 단계와 개념에 대해 자세히 알아보려면 다음 리소스를 참조하세요.
스트림 프로세서를 생성, 시작, 중지 및 모니터링하는 방법을 학습하려면 스트림 프로세서 개발을(를) 참조하세요.
Atlas Stream Processing이 지원하는 집계 단계에 대해 자세히 학습하려면 집계 파이프라인 단계를 참조하세요.
스트림 프로세서가 읽을 수 있는 소스에 대해 자세히 알아보려면
$source단계(Stream Processing)를 참조하세요.창 단계에 대해 학습하려면, 스트림 프로세서 Windows를 참조하십시오.
온디맨드 구체화된 보기에 대한 자세한 내용은 온디맨드 구체화 보기를 참조하세요.