MongoDB C++ 드라이버 심화: 집계 파이프라인($lookup), 인덱스 전략, Read Preference
C++에서 MongoDB 쓰기(#52-3)에서 mongocxx 설치와 CRUD, 기본적인 집계와 인덱스를 다뤘습니다. 이 글은 그다음 단계로, 이벤트 로그처럼 문서가 많이 쌓이는 컬렉션에서 통계를 뽑고, 그 조회를 인덱스로 빠르게 만들고, 레플리카셋에서 읽기를 어디로 보낼지 정하는 방법을 다룹니다. 예제는 mongocxx 4.x와 C++17을 기준으로 하며, 아래 using 선언을 전제로 합니다.
#include <bsoncxx/builder/basic/array.hpp>
#include <bsoncxx/builder/basic/document.hpp>
#include <bsoncxx/json.hpp>
#include <bsoncxx/types.hpp>
#include <mongocxx/client.hpp>
#include <mongocxx/instance.hpp>
#include <mongocxx/pipeline.hpp>
#include <mongocxx/options/aggregate.hpp>
#include <mongocxx/options/find.hpp>
#include <mongocxx/options/index.hpp>
#include <mongocxx/read_preference.hpp>
#include <mongocxx/uri.hpp>
#include <mongocxx/write_concern.hpp>
#include <chrono>
#include <iostream>
using bsoncxx::builder::basic::kvp;
using bsoncxx::builder::basic::make_array;
using bsoncxx::builder::basic::make_document;
집계 파이프라인
find는 필터, 정렬, 프로젝션까지만 할 수 있습니다. 그룹별 합계나 평균, 다른 컬렉션과의 결합처럼 SQL의 GROUP BY나 JOIN에 해당하는 작업은 aggregate로 합니다. 파이프라인은 스테이지의 배열이고, 각 스테이지의 출력이 다음 스테이지의 입력이 됩니다. mongocxx에서는 mongocxx::pipeline 객체에 스테이지를 이어 붙입니다.
시간대별 평균 응답 시간
void hourly_latency(mongocxx::collection& events) {
using namespace std::chrono;
auto since = bsoncxx::types::b_date{system_clock::now() - hours(24)};
mongocxx::pipeline p;
p.match(make_document(
kvp("type", "api_call"),
kvp("createdAt", make_document(kvp("$gte", since)))))
.group(make_document(
kvp("_id", make_document(kvp("$dateTrunc", make_document(
kvp("date", "$createdAt"), kvp("unit", "hour"), kvp("timezone", "Asia/Seoul"))))),
kvp("avgMs", make_document(kvp("$avg", "$latencyMs"))),
kvp("maxMs", make_document(kvp("$max", "$latencyMs"))),
kvp("calls", make_document(kvp("$sum", 1)))))
.sort(make_document(kvp("_id", 1)));
for (auto&& doc : events.aggregate(p)) {
std::cout << bsoncxx::to_json(doc) << "\n";
}
}
$dateTrunc(MongoDB 5.0 이상)는 날짜를 시간 단위로 잘라 Date 타입 그대로 돌려주므로, $dateToString으로 문자열을 만드는 것보다 정렬과 후처리가 편합니다. 시간대를 지정하지 않으면 UTC 기준으로 자르므로 한국 시간 기준 통계라면 timezone을 꼭 넘깁니다. 첫 스테이지인 $match는 {type: 1, createdAt: 1} 같은 인덱스를 쓸 수 있습니다. $group 뒤에서는 인덱스를 쓸 수 없으므로, 필터 조건은 항상 파이프라인 앞쪽에 둡니다. 서버의 최적화기가 일부 스테이지 순서를 자동으로 바꿔 주기도 하지만, 처음부터 그렇게 쓰는 편이 의도가 분명합니다.
$lookup과 타입 불일치
void orders_with_user(mongocxx::collection& orders) {
mongocxx::pipeline p;
p.lookup(make_document(
kvp("from", "users"),
kvp("localField", "userId"),
kvp("foreignField", "_id"),
kvp("as", "user")));
p.unwind(make_document(
kvp("path", "$user"),
kvp("preserveNullAndEmptyArrays", true))); // 사용자가 없는 주문도 유지
p.project(make_document(
kvp("amount", 1),
kvp("userName", "$user.name")));
for (auto&& doc : orders.aggregate(p)) {
std::cout << bsoncxx::to_json(doc) << "\n";
}
}
$lookup의 결과가 늘 빈 배열이라면 거의 항상 두 필드의 타입이 다른 경우입니다. 주문 문서의 userId는 문자열 "64f..."인데 users._id는 ObjectId라면 값이 같아 보여도 매칭되지 않습니다. 가장 좋은 해결책은 저장할 때부터 타입을 통일하는 것이고, 이미 섞여 있다면 lookup 전에 변환합니다.
p.add_fields(make_document(
kvp("userOid", make_document(kvp("$toObjectId", "$userId")))));
p.lookup(make_document(
kvp("from", "users"),
kvp("localField", "userOid"),
kvp("foreignField", "_id"),
kvp("as", "user")));
$toObjectId는 24자리 16진수가 아닌 문자열을 만나면 오류를 내므로, 형식이 섞여 있다면 $convert에 onError를 지정해 씁니다. $lookup은 입력 문서마다 from 컬렉션을 조회하므로 foreignField에 인덱스가 없으면 입력 문서 수만큼 컬렉션 스캔이 반복됩니다. _id는 기본 인덱스가 있지만, 다른 필드로 결합한다면 인덱스를 먼저 만들어야 합니다.
$facet: 같은 입력으로 여러 통계 만들기
void dashboard(mongocxx::collection& orders) {
mongocxx::pipeline p;
p.match(make_document(kvp("status", "paid")));
p.facet(make_document(
kvp("byDay", make_array(
make_document(kvp("$group", make_document(
kvp("_id", make_document(kvp("$dateTrunc", make_document(
kvp("date", "$createdAt"), kvp("unit", "day"))))),
kvp("total", make_document(kvp("$sum", "$amount")))))),
make_document(kvp("$sort", make_document(kvp("_id", 1)))))),
kvp("topUsers", make_array(
make_document(kvp("$group", make_document(
kvp("_id", "$userId"),
kvp("orders", make_document(kvp("$sum", 1)))))),
make_document(kvp("$sort", make_document(kvp("orders", -1)))),
make_document(kvp("$limit", 10))))));
for (auto&& doc : orders.aggregate(p)) {
std::cout << bsoncxx::to_json(doc) << "\n"; // {byDay: [...], topUsers: [...]} 문서 하나
}
}
$facet은 앞 스테이지의 결과를 여러 하위 파이프라인에 똑같이 흘려보내고, 각 결과를 배열 필드로 담은 문서 하나를 돌려줍니다. 하위 파이프라인이 병렬로 실행되는 것은 아니며, 이득은 입력을 한 번만 읽고 왕복을 한 번으로 줄이는 데 있습니다. 주의할 점이 두 가지 있습니다. 하위 파이프라인 안에서는 인덱스를 쓸 수 없으므로, 공통 필터는 위 코드처럼 $facet 앞의 $match에 둬야 합니다. 그리고 결과가 문서 하나이므로 BSON 문서 크기 제한인 16MB를 넘을 수 없습니다. 하위 결과가 클 수 있다면 $limit으로 잘라야 합니다.
$merge로 통계를 컬렉션에 저장
대시보드가 열릴 때마다 큰 집계를 돌리는 대신, 주기적으로 집계해 결과를 별도 컬렉션에 저장해 두고 그것을 읽는 방법이 흔히 쓰입니다.
void refresh_daily_stats(mongocxx::database db) {
using namespace std::chrono;
mongocxx::pipeline p;
p.match(make_document(kvp("createdAt", make_document(
kvp("$gte", bsoncxx::types::b_date{system_clock::now() - hours(48)})))));
p.group(make_document(
kvp("_id", make_document(kvp("$dateTrunc", make_document(
kvp("date", "$createdAt"), kvp("unit", "day"), kvp("timezone", "Asia/Seoul"))))),
kvp("events", make_document(kvp("$sum", 1)))));
p.merge(make_document(
kvp("into", "daily_stats"),
kvp("on", "_id"),
kvp("whenMatched", "replace"),
kvp("whenNotMatched", "insert")));
// $merge는 결과를 반환하지 않지만, 커서를 순회해야 파이프라인이 실행됩니다.
auto cursor = db["events"].aggregate(p);
for (auto&& _ : cursor) { (void)_; }
}
$merge(MongoDB 4.2 이상)는 지정한 키로 기존 문서와 맞춰 갱신하거나 삽입하므로, 최근 이틀 치만 다시 집계해도 오래된 날짜의 통계는 그대로 남습니다. $out은 대상 컬렉션 전체를 결과로 교체하므로 증분 갱신에는 맞지 않습니다. mongocxx의 aggregate는 커서를 돌려주기만 하고, 서버에 명령을 보내는 것은 커서를 처음 순회할 때입니다. 위처럼 빈 루프라도 돌려야 실제로 실행됩니다.
메모리 한도와 allowDiskUse
$group, $sort, $bucket 같은 스테이지는 스테이지별로 100MB 메모리 한도가 있고, 넘으면 Exceeded memory limit 오류가 납니다. mongocxx::options::aggregate의 allow_disk_use(true)를 켜면 임시 파일을 사용해 계속 진행합니다. MongoDB 6.0부터는 서버 파라미터 allowDiskUseByDefault가 기본으로 켜져 있어 명시하지 않아도 디스크를 쓰지만, 명시하면 이전 버전 서버에서도 같은 동작을 기대할 수 있습니다. $graphLookup처럼 이 옵션의 적용 범위가 다른 스테이지도 있으므로, 사용하는 서버 버전의 문서에서 해당 스테이지의 제한을 확인하는 것이 좋습니다. 디스크를 쓰면 느려지므로 먼저 $match와 $project로 입력을 줄이는 것이 우선입니다.
인덱스 설계
인덱스가 없는 조회는 컬렉션 전체를 읽는 COLLSCAN이 되고, 문서 수에 비례해 느려집니다. 인덱스를 만들면 B-tree를 따라 필요한 범위만 읽는 IXSCAN이 됩니다. 대신 인덱스마다 쓰기 시 갱신 비용과 메모리(작업 집합)를 차지하므로, 실제 조회 패턴에 맞춘 소수의 인덱스를 두는 것이 목표입니다.
복합 인덱스의 필드 순서
복합 인덱스는 앞쪽 필드부터 차례로 정렬된 구조라서, 필드 순서가 어떤 조회에 쓰일 수 있는지를 결정합니다. 일반적인 원칙은 ESR(Equality, Sort, Range)입니다. 등호 조건 필드를 먼저, 정렬 필드를 다음에, 범위 조건 필드를 마지막에 둡니다.
void create_event_indexes(mongocxx::collection& events) {
// 조회: userId == ? 이고 createdAt 역순 정렬, 최근 N개
events.create_index(make_document(kvp("userId", 1), kvp("createdAt", -1)));
// 조회: type == ? 이고 createdAt 범위
events.create_index(make_document(kvp("type", 1), kvp("createdAt", 1)));
}
{userId: 1, createdAt: -1} 인덱스는 userId만으로 조회할 때도 쓰이지만(접두사), createdAt만으로 조회할 때는 쓰이지 않습니다. 범위 조건을 정렬 필드보다 앞에 두면, 범위에 걸리는 여러 구간을 합친 뒤 메모리에서 다시 정렬해야 해서(SORT 스테이지) 인덱스의 정렬 순서를 쓸 수 없습니다.
부분 인덱스와 유니크 제약
void create_partial_unique(mongocxx::collection& users) {
mongocxx::options::index opts{};
opts.unique(true);
// email 필드가 있는 문서에만 유니크 제약을 건다
opts.partial_filter_expression(make_document(
kvp("email", make_document(kvp("$exists", true)))));
users.create_index(make_document(kvp("email", 1)), opts);
}
일반 유니크 인덱스는 필드가 없는 문서도 null 값으로 인덱싱하므로, email이 없는 문서가 두 개만 되어도 중복 키 오류가 납니다. 부분 인덱스로 대상 문서를 제한하면 이 문제를 피할 수 있고, 인덱스 크기도 줄어듭니다. 단, 조회 조건이 부분 인덱스의 필터 조건을 포함해야 그 인덱스가 선택됩니다.
텍스트 인덱스
void text_search(mongocxx::collection& articles) {
mongocxx::options::index idx_opts{};
idx_opts.default_language("english");
articles.create_index(make_document(kvp("title", "text"), kvp("content", "text")), idx_opts);
mongocxx::options::find opts{};
// 옵션에는 value를 넘겨야 합니다. 임시 value의 view를 넘기면 dangling이 됩니다.
opts.projection(make_document(kvp("title", 1),
kvp("score", make_document(kvp("$meta", "textScore")))));
opts.sort(make_document(kvp("score", make_document(kvp("$meta", "textScore")))));
auto filter = make_document(kvp("$text", make_document(kvp("$search", "mongodb driver"))));
for (auto&& doc : articles.find(filter.view(), opts)) {
std::cout << bsoncxx::to_json(doc) << "\n";
}
}
컬렉션당 텍스트 인덱스는 하나만 만들 수 있으므로 검색 대상 필드를 한 인덱스에 모두 넣습니다. 텍스트 인덱스는 언어별 어간 추출과 불용어 처리를 하는데, 지원 언어 목록에 한국어는 없습니다. 한국어 본문을 검색해야 한다면 default_language("none")으로 단순 토큰 매칭을 하거나, 형태소 분석이 필요한 수준이라면 Atlas Search나 Elasticsearch 같은 별도 검색 엔진을 쓰는 것이 현실적입니다.
인덱스 생성 시 주의
create_index는 같은 키와 같은 옵션의 인덱스가 이미 있으면 아무 일도 하지 않고 성공합니다. 그래서 앱 기동 시 필요한 인덱스를 매번 생성 요청해도 됩니다. 같은 키에 다른 옵션(예: unique 여부)으로 만들려고 하면 IndexOptionsConflict(코드 85)나 IndexKeySpecsConflict(코드 86) 오류가 납니다. 이 오류를 잡아서 무시하면 의도한 옵션이 적용되지 않은 채로 운영하게 되므로, 기존 인덱스를 지우고 다시 만들어야 하는 상황인지 확인하는 신호로 다뤄야 합니다. 이미 데이터가 많은 컬렉션에서 유니크 인덱스를 만들 때 기존 데이터에 중복이 있으면 생성이 실패하므로, 먼저 집계로 중복을 찾아 정리합니다.
explain으로 실행 계획 확인
mongocxx의 find 옵션에는 explain 설정이 없으므로 explain 명령을 직접 보냅니다.
void explain_find(mongocxx::database db) {
auto plan = db.run_command(make_document(
kvp("explain", make_document(
kvp("find", "events"),
kvp("filter", make_document(kvp("userId", "user123"))),
kvp("sort", make_document(kvp("createdAt", -1))))),
kvp("verbosity", "executionStats")));
std::cout << bsoncxx::to_json(plan) << "\n";
}
결과에서 볼 부분은 세 가지입니다. queryPlanner.winningPlan에 IXSCAN이 있는지(없고 COLLSCAN이면 인덱스를 못 쓴 것), SORT 스테이지가 있는지(있으면 인덱스 순서로 정렬하지 못하고 메모리에서 정렬한 것), 그리고 executionStats의 totalDocsExamined가 nReturned보다 훨씬 큰지입니다. 반환한 문서 수에 비해 읽은 문서 수가 많다면 인덱스가 조건을 충분히 좁히지 못하고 있다는 뜻입니다.
레플리카셋: Read Preference와 Write Concern
레플리카셋은 쓰기를 받는 Primary 하나와 Primary의 oplog를 복제하는 Secondary들로 구성됩니다. URI에 replicaSet 이름과 호스트 목록을 주면 드라이버가 구성 전체를 파악하고, Primary가 바뀌면 새 Primary를 찾아 쓰기를 보냅니다.
mongocxx::uri uri("mongodb://host1:27017,host2:27017,host3:27017/?replicaSet=rs0");
replicaSet 옵션 없이 호스트 하나만 적으면 드라이버는 그 서버에 직접 연결합니다(direct connection). 그 서버가 지금 Primary라면 문제없이 동작하다가, 장애 조치로 Secondary가 되는 순간 쓰기가 NotWritablePrimary 오류로 실패하기 시작합니다. 평소에는 드러나지 않다가 장애 조치 때 터지는 설정 실수이므로, 레플리카셋에는 항상 replicaSet 이름을 넣습니다.
Read Preference
| 모드 | 동작 |
|---|---|
k_primary | Primary에서만 읽음 (기본값). Primary가 없으면 실패 |
k_primary_preferred | Primary 우선, 없으면 Secondary |
k_secondary | Secondary에서만 읽음. Secondary가 없으면 실패 |
k_secondary_preferred | Secondary 우선, 없으면 Primary |
k_nearest | 응답 지연이 가장 짧은 노드(일정 범위 안에서 무작위) |
mongocxx::cursor read_report(mongocxx::collection& events) {
mongocxx::read_preference rp{};
rp.mode(mongocxx::read_preference::read_mode::k_secondary_preferred);
// 복제 지연이 90초를 넘는 Secondary는 선택하지 않음
rp.max_staleness(std::chrono::seconds{90});
mongocxx::options::aggregate opts{};
opts.read_preference(rp);
opts.allow_disk_use(true);
mongocxx::pipeline p;
p.group(make_document(kvp("_id", "$type"), kvp("n", make_document(kvp("$sum", 1)))));
return events.aggregate(p, opts);
}
Secondary 읽기는 복제 지연만큼 오래된 데이터를 볼 수 있으므로, 리포트나 분석처럼 약간의 지연이 허용되는 무거운 읽기에 적합합니다. 사용자가 방금 저장한 데이터를 바로 보여 주는 화면이라면 Primary에서 읽어야 합니다. max_staleness는 90초 이상으로만 설정할 수 있습니다. 단독 서버에 연결했다면 read preference는 무시되고 그 서버에서 읽으므로 오류는 나지 않지만 분산 효과도 없습니다. Secondary로 읽기를 분산한다고 쓰기 부하가 줄지는 않으며, Secondary도 같은 쓰기를 모두 복제 적용하고 있다는 점을 염두에 둬야 합니다.
Write Concern
void insert_durable(mongocxx::collection& payments, bsoncxx::document::view doc) {
mongocxx::write_concern wc{};
wc.acknowledge_level(mongocxx::write_concern::level::k_majority);
wc.journal(true);
wc.timeout(std::chrono::milliseconds{5000});
mongocxx::options::insert opts{};
opts.write_concern(wc);
payments.insert_one(doc, opts);
}
w: majority는 과반수 노드가 쓰기를 기록한 뒤 응답하므로, Primary가 장애로 바뀌어도 응답을 받은 쓰기는 롤백되지 않습니다. MongoDB 5.0부터는 대부분의 구성에서 기본 write concern이 이미 majority입니다. timeout(wtimeout)을 넘기면 오류가 나지만 쓰기가 취소되는 것은 아닙니다. Primary에는 이미 기록되었고 복제 확인만 시간 안에 받지 못한 것이므로, 오류를 받았다고 같은 쓰기를 그대로 다시 보내면 중복될 수 있습니다. 재시도할 쓰기는 upsert나 고유 키로 멱등하게 만들어 두는 것이 안전합니다.
자주 만나는 오류
Exceeded memory limit 오류는 $group이나 $sort가 100MB 한도를 넘은 경우입니다. 앞쪽 $match/$project로 입력을 줄이고, 필요하면 allow_disk_use(true)를 켭니다.
$lookup 결과가 빈 배열이라면 localField와 foreignField의 타입(문자열과 ObjectId 등)이 같은지 먼저 확인합니다. 앞의 $toObjectId 변환 예제를 참고합니다.
read preference를 secondary로 설정했는데 Primary에서 읽힌다면, 옵션 객체를 find/aggregate 호출에 실제로 넘겼는지, 연결이 레플리카셋 모드인지(URI에 replicaSet이 있는지) 확인합니다. 컬렉션 수준에서 기본값을 정하고 싶다면 collection.read_preference(rp)로 설정할 수 있습니다.
옵션에 정렬이나 프로젝션을 make_document(...).view() 형태로 넘기면, 임시 value가 그 문장 끝에서 파괴되어 실제 호출 시점에는 해제된 메모리를 가리킵니다. 항상 value를 넘기거나, value를 변수에 담아 호출이 끝날 때까지 살려 둬야 합니다.
다음 글: C++ PostgreSQL 드라이버(#52-4) 이전 글: C++에서 MongoDB 쓰기
참고 자료
- MongoDB Aggregation Pipeline
- MongoDB Index Strategies
- MongoDB Replica Set
- mongocxx API 레퍼런스
- MongoDB C++ Driver 공식 문서