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 스테이지가 사용되는 경우, 참조 컬렉션의 나중 변경으로 인해 뷰가 이미 쓰기 (write)한 문서는 업데이트되지 않습니다.

다음 절차는 구매 방식에 따라 완료된 판매 건수를 유지하는 스트리밍 구체화된 뷰를 만듭니다. 스트림 프로세서는 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 작업 공간 이 있어야 합니다. 연결을 추가하려면 연결 관리를 참조하세요.

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

1

이 절차에서는 sample_supplies 데이터 세트의 sales 컬렉션을 사용합니다. 샘플 데이터를 로드하는 방법은 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 컬렉션 변경 스트림을 참조하십시오.

1
  1. Atlas UI에서 Atlas 프로젝트의 Stream Processing 페이지로 Go.

  2. 소스 클러스터에 대한 Atlas 연결이 있는 Stream Processing 워크스페이스 창의 Manage 을 클릭합니다.

2
  1. Create stream processor를 클릭합니다.

  2. Visual Builder 을(를) 선택합니다.

3
  1. Source 필드에서 Connection 드롭다운 목록에서 소스 클러스터에 대한 Atlas 연결을 선택합니다.

  2. JSON 텍스트 상자에서 $source 단계를 구성하여 사전 및 게시 이미지로 sales 컬렉션을 읽습니다.

{
"$source": {
"connectionName": "<connection-name>",
"db": "sample_supplies",
"coll": "sales",
"config": {
"fullDocument": "required",
"fullDocumentBeforeChange": "required"
}
}
}
4
  1. Start building your pipeline 창에서 + Custom stage을 클릭합니다.

  2. 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"
]
}
}
}
5
  1. +을 클릭한 다음 Custom stage을 선택합니다.

  2. JSON 텍스트 상자에서 델타가 0인 이벤트를 제거하는 $match 단계를 추가합니다.

{
"$match": { "_delta": { "$ne": 0 } }
}
6
  1. +을 클릭한 다음 Custom stage을 선택합니다.

  2. 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 단계가 창 단계 내에서 실행되어야 합니다.

7
  1. Sink 필드에서 Connection 드롭다운 목록에서 사용자의 Atlas 연결을 선택합니다.

  2. 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 최고 수위는 재생된 창이 이중 계산되는 것을 방지합니다.

8
  1. Stream processor name 필드에 sales_stats_sp를 입력합니다.

  2. Stream Processing의 계층을 선택합니다. 워크로드에 맞는 계층을 선택하려면 Atlas Stream Processing 계층 선택 가이드를 참조하세요.

9

Create stream processor를 클릭합니다.

10

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

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

  2. $addFields 단계는 operationType 에서 $switch 를 사용하여 _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. $addFields 다음에 $unwind 단계와 $set 단계를 추가합니다. $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"
}
}

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