AI 에이전트의 경우: 문서 인덱스는 https://www.mongodb.com/ko-kr/docs/llms.txt에서 사용할 수 있으며, 모든 페이지의 마크다운 버전은 어떤 URL 경로에 .md를 추가하여 사용할 수 있습니다.
Docs Menu

스트리밍 구체화된 뷰 빌드

스트리밍 구체화된 뷰는 Atlas Stream Processing 스트림 프로세서가 현재 상태를 유지하는 컬렉션입니다. 이를 위해 스트림 프로세서는 소스 컬렉션의 각 변경 사항을 읽고, 변경의 영향을 계산하고, 결과를 뷰에 적용합니다. 이 뷰는 수동 또는 예정된 새로 고침 없이 소스 데이터의 현재 상태를 반영합니다.

스트리밍 구체화된 뷰는 데이터 소스가 변경되는 경우보다 계산된 결과를 읽는 워크로드(예: 대시보드, 누계 합계, 미결 항목 개수) 에 적합합니다. 데이터가 도착하면 뷰가 업데이트되므로 읽기는 지연 시간이 즐어들어 현재 결과를 반환합니다.

MongoDB는 구체화된 뷰에 대한 두 가지 접근 방식을 지원합니다.

  • 온디맨드 구체화된 뷰는 수동으로 또는 예정 에 따라 실행 집계 파이프라인 의 결과를 저장합니다.

  • 스트리밍 구체화된 뷰는 연속적으로 실행되는 Atlas Stream Processing 스트리밍 프로세서의 결과를 저장합니다.

두 개 방법 모두 계산된 결과를 디스크에 저장하고 보기에서 직접 읽기를 제공합니다. 보기가 업데이트되는 방법과 시기가 다릅니다.

특성
온디맨드 구체화된 뷰
스트리밍 구체화된 뷰

trigger 업데이트

수동 또는 예정

변경 중심의 연속적

지연 시간

분을 일수로 변환

초 미만에서 수 초 단위

데이터 신선도

특정 점 스냅샷

영구적으로 동기화됨

컴퓨트 모델

전체 결과 다시 계산

증분 효과 계산

Best fit

배치 보고, 주기적 집계

실시간 대시보드, 운영 분석

스트리밍 구체화된 뷰를 유지하는 것은 표준 Atlas Stream Processing 집계 단계로 구성하는 패턴입니다. 이 패턴을 구현하는 스트리밍 프로세서에는 다음과 같은 특성이 있을 수 있습니다.

  • 변경 스트림 소스를 읽습니다. $source 단계는 fullDocumentfullDocumentBeforeChangerequired (으)로 설정된 소스 컬렉션에서 읽어오므로 파이프라인은 변경 전후 각 문서의 상태를 비교할 수 있습니다.

  • 이벤트 당 부호 있는 델타를 계산합니다. $addFields 단계에서는 $switch 표현식 사용하여 변경 사항이 계산된 결과에 미치는 영향에 따라 각 삽입, 업데이트 또는 삭제 에 양수 또는 음수 값을 할당할 수 있습니다.

  • 결과를 창 으로 그룹화합니다. 스트림 프로세서는 제한 없는 스트림 $group 에서 작동하므로 모든 단계는 창 단계 내에서 실행 되어야 합니다. 윈도우 는 $group 창 간격 동안 각 키의 델타를 합산합니다. 창 간격은 뷰의 새로 고침 간격으로도 작용하므로 뷰의 최신 상태를 결정합니다. 설정하다 수 있는 가장 작은 간격은 1밀리초입니다.

  • 결과를 추가로 적용합니다. $merge 단계에서는 whenMatched 각 창의 결과를 대체하는 대신 뷰의 실행 합계에 추가하는 파이프라인 사용할 수 있습니다.

  • 싱크 단계에서 종료됩니다. 스트림 프로세서 파이프라인 싱크 단계에서 끝나야 합니다. 를 사용하여 Atlas 컬렉션 에 쓰기 (write) $merge .

증분 집계만 스트리밍 구체화된 뷰로 변환됩니다. 전체 컬렉션 스캔이 필요한 집계는 자격이 없습니다.

스트리밍 구체화된 뷰는 하위 소비자에 영향을 미치는 방식으로 배치 집계와 다르게 동작합니다.

  • 보기는 0부터 시작됩니다. 기본적으로 프로세서는 기존 문서를 읽지 않으므로 프로세서를 시작하기 전에 보기를 시드하거나 $source 단계에서 initialSync 을(를) 활성화합니다.

  • 소스 변경만 업데이트를 유도합니다. 파이프라인 단계를 사용하는 경우 참조 $lookup 컬렉션 나중에 변경해도 뷰가 이미 작성한 문서는 업데이트 되지 않습니다.

다음 튜토리얼은 구매 방식에 따라 완료된 판매 건수를 유지하는 스트리밍 구체화된 뷰를 만듭니다. 스트림 프로세서는 sample_supplies.sales 컬렉션의 변경 스트림을 읽고 sample_supplies.sales_by_channel 컬렉션에 쓰기 (write)를 수행합니다. [[ ## completed ##]]

참고

스트림 프로세서는 또한 Apache Kafka 주제 에서 읽고 AWS S3의 Apache ice버그 테이블에 쓰기 (write) 수 있습니다. 자세한 학습 은 Apache Kafka 브로커$iceberg 애그리게이션 단계를 참조하세요.

Stream Processor를 만들기 전에 소스 데이터가 있는 클러스터에 대한 Atlas 연결이 있는 Stream Processing 작업 공간 이 있어야 합니다. 연결을 추가하려면 연결 관리를 참조하세요.

이 섹션의 명령을 클러스터에 대해 실행합니다. 연결하려면 mongosh 를 통해 클러스터에 연결을 참조하세요.

다음 단계를 완료하여 소스 컬렉션을 준비하고 보기를 시딩합니다.

1

이 절차에서는 sales sample_supplies 데이터 세트의 컬렉션 사용합니다.샘플 데이터를 로드하는 방법을 학습 Atlas 배포로 샘플 데이터 가져오기를 참조하세요.

2

status 필드를 추가하여 스트림 프로세서가 완료된 및 반환된 판매를 감지할 수 있도록 합니다. 클러스터에 대해 다음 명령을 실행합니다.

db.sales.updateMany(
{ status: { $exists: false } },
{ $set: { status: "completed" } }
)
{
acknowledged: true,
insertedId: null,
matchedCount: 5000,
modifiedCount: 5000,
upsertedCount: 0
}
3

스트림 프로세서가 각 변경 사항의 영향을 계산할 수 있도록 원본 컬렉션에서 사전 및 게시 이미지를 활성화합니다. 클러스터에 대해 다음 명령을 실행합니다.

db.getSiblingDB("sample_supplies").runCommand({
collMod: "sales",
changeStreamPreAndPostImages: { enabled: true }
})
{ ok: 1, ... }
4

클러스터에 대해 다음 배치 집계를 실행하여 현재 카운트로 뷰를 채우세요.

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 컬렉션 변경 스트림을 참조하십시오.

Atlas UI 또는 를 선택하여 스트림 프로세서를 mongosh 만듭니다.

이 섹션의 명령을 클러스터에 대해 실행합니다. 연결하려면 mongosh 를 통해 클러스터에 연결을 참조하세요.

프로세서를 시작하면 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 로 티켓의 우선 순위에 따라 카운트를 조정합니다. 집계에는 다섯 단계가 있습니다.

  1. $source 단계는 사전 이미지와 게시 이미지로 support_tickets 변경 스트림을 읽습니다.

  2. 단계에서는 $addFields $switch에서 operationType 을 사용하여 _delta 를 계산하고 를 _priority 그룹 키로 추출합니다.

  3. $match 단계에서는 델타가 0인 이벤트를 삭제합니다.

  4. $tumblingWindow 단계는 각 1초 창 내에서 우선 순위별 델타를 합산합니다.

  5. $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 배열로 대체하고 배열을 조정당 하나의 문서로 팬아웃합니다. 단일 키 파이프라인을 다음과 같이 수정합니다.

  1. $addFields $switch_adjustments 분기가 영향을 받는 그룹 키당 하나의 요소가 포함된 배열 반환하도록 단계를 대체합니다. 에스컬레이션은 두 개의 요소를 반환합니다.

  2. $unwind $set $addFields뒤에 $unwind 단계와 $match 단계를 추가합니다. 단계는 조정당 각 이벤트 하나의 문서 로 분할하고 단계를 대체하는 빈 배열을 삭제합니다.$set 단계는 _adjustments 필드를 최상위 수준 _priority 및 로 승격하여 _delta 나머지 단계는 변경되지 않고 작동하도록 합니다.

단일 키 파이프라인의 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"
}
}

이 가이드의 단계와 개념에 대해 자세히 알아보려면 다음 리소스를 참조하세요.