БСсплатная ΠΊΠΎΠ½ΡΡƒΠ»ΡŒΡ‚Π°Ρ†ΠΈΡ

ΠžΠΏΠΈΡˆΠΈΡ‚Π΅ Π·Π°Π΄Π°Ρ‡Ρƒ ΠΈΠ»ΠΈ Ρ†Π΅Π»ΡŒ. ΠžΡ‚Π²Π΅Ρ‡Ρƒ с практичСским ΡΠ»Π΅Π΄ΡƒΡŽΡ‰ΠΈΠΌ шагом β€” бСсплатно, Π±Π΅Π· ΠΎΠ±ΡΠ·Π°Ρ‚Π΅Π»ΡŒΡΡ‚Π².

Или Π²Ρ‹Π±Π΅Ρ€ΠΈΡ‚Π΅ врСмя Π² Calendly

Change Streams Π² экосистСмС MongoDB

MongoDB Change Streams β€” API для подписки Π½Π° ΠΏΠΎΡ‚ΠΎΠΊ ΠΈΠ·ΠΌΠ΅Π½Π΅Π½ΠΈΠΉ Π² коллСкциях, Π±Π°Π·Π°Ρ… ΠΈΠ»ΠΈ всём кластСрС (начиная с replica set ΠΈΠ»ΠΈ sharded cluster). ВмСсто polling ΠΏΠΎ updatedAt ΠΈΠ»ΠΈ tailing oplog Π²Ρ€ΡƒΡ‡Π½ΡƒΡŽ ΠΏΡ€ΠΈΠ»ΠΎΠΆΠ΅Π½ΠΈΠ΅ ΠΏΠΎΠ»ΡƒΡ‡Π°Π΅Ρ‚ упорядочСнныС события insert, update, replace, delete с ΠΏΠΎΠ»Π½Ρ‹ΠΌ ΠΈΠ»ΠΈ delta Π΄ΠΎΠΊΡƒΠΌΠ΅Π½Ρ‚ΠΎΠΌ послС измСнСния.

Π­Ρ‚ΠΎ основа для real-time Π΄Π°ΡˆΠ±ΠΎΡ€Π΄ΠΎΠ², ΠΈΠ½Π²Π°Π»ΠΈΠ΄Π°Ρ†ΠΈΠΈ кэша, синхронизации search index (Elasticsearch, Atlas Search), event-driven микросСрвисов, audit log ΠΈ push-ΡƒΠ²Π΅Π΄ΠΎΠΌΠ»Π΅Π½ΠΈΠΉ ΠΏΠΎΠ»ΡŒΠ·ΠΎΠ²Π°Ρ‚Π΅Π»ΡΠΌ. Change Streams ΠΎΠΏΠΈΡ€Π°ΡŽΡ‚ΡΡ Π½Π° WiredTiger oplog ΠΈ inherit Π΅Π³ΠΎ guarantees ordering Π² ΠΏΡ€Π΅Π΄Π΅Π»Π°Ρ… ΠΎΠ΄Π½ΠΎΠ³ΠΎ shard key для sharded collections.

Как устроСн ΠΏΠΎΡ‚ΠΎΠΊ событий

ΠšΠ»ΠΈΠ΅Π½Ρ‚ ΠΎΡ‚ΠΊΡ€Ρ‹Π²Π°Π΅Ρ‚ cursor Ρ‡Π΅Ρ€Π΅Π· collection.watch() ΠΈΠ»ΠΈ db.watch() с pipeline Π°Π³Ρ€Π΅Π³Π°Ρ†ΠΈΠΈ для Ρ„ΠΈΠ»ΡŒΡ‚Ρ€Π°Ρ†ΠΈΠΈ. КаТдоС событиС содСрТит:

  • operationType β€” insert, update, replace, delete, invalidate ΠΈ Π΄Ρ€.
  • fullDocument β€” ΠΏΡ€ΠΈ updateLookup ΠΈΠ»ΠΈ ΠΏΡ€ΠΈ insert/replace;
  • documentKey β€” _id ΠΈΠ·ΠΌΠ΅Π½Ρ‘Π½Π½ΠΎΠ³ΠΎ Π΄ΠΎΠΊΡƒΠΌΠ΅Π½Ρ‚Π°;
  • clusterTime β€” логичСская ΠΌΠ΅Ρ‚ΠΊΠ° для resume;
  • updateDescription β€” changedFields ΠΈ removedFields для update.

Resume token позволяСт ΠΏΠ΅Ρ€Π΅ΠΏΠΎΠ΄ΠΊΠ»ΡŽΡ‡ΠΈΡ‚ΡŒΡΡ послС ΠΎΠ±Ρ€Ρ‹Π²Π° Π±Π΅Π· пропуска ΠΈΠ»ΠΈ дублирования (ΠΏΡ€ΠΈ ΠΊΠΎΡ€Ρ€Π΅ΠΊΡ‚Π½ΠΎΠΉ ΠΎΠ±Ρ€Π°Π±ΠΎΡ‚ΠΊΠ΅ idempotency Π½Π° consumer). Π‘Π΅Π· сохранСния token consumer послС restart Π½Π°Ρ‡Π½Ρ‘Ρ‚ с Β«nowΒ» ΠΈ потСряСт ΠΏΡ€ΠΎΠΌΠ΅ΠΆΡƒΡ‚ΠΎΠΊ.

Pipeline-Ρ„ΠΈΠ»ΡŒΡ‚Ρ€Π°Ρ†ΠΈΡ Π½Π° сСрвСрС

$match Π½Π° operationType, поля fullDocument.status ΠΈΠ»ΠΈ tenantId сниТаСт Ρ‚Ρ€Π°Ρ„ΠΈΠΊ Π½Π° consumer. Π€ΠΈΠ»ΡŒΡ‚Ρ€ΡƒΠΉΡ‚Π΅ ΠΊΠ°ΠΊ ΠΌΠΎΠΆΠ½ΠΎ Π±Π»ΠΈΠΆΠ΅ ΠΊ источнику β€” Π½Π΅ Ρ‚Π°Ρ‰ΠΈΡ‚Π΅ вСсь oplog Π² ΠΏΡ€ΠΈΠ»ΠΎΠΆΠ΅Π½ΠΈΠ΅ Ρ€Π°Π΄ΠΈ ΠΎΠ΄Π½ΠΎΠ³ΠΎ поля.

АрхитСктура real-time прилоТСния

Випичная схСма: MongoDB β†’ Change Stream consumer (Node, Python, Go worker) β†’ message bus (Kafka, Redis Streams, NATS) ΠΈΠ»ΠΈ WebSocket gateway β†’ ΠΊΠ»ΠΈΠ΅Π½Ρ‚Ρ‹. ΠΠ»ΡŒΡ‚Π΅Ρ€Π½Π°Ρ‚ΠΈΠ²Π° β€” serverless trigger Ρ‡Π΅Ρ€Π΅Π· Atlas Trigger ΠΈΠ»ΠΈ Change Stream Handler Π² Ρ‚ΠΎΠΌ ΠΆΠ΅ процСссС API для простых кСйсов.

WebSocket ΠΈ SSE

API-сСрвСр Π΄Π΅Ρ€ΠΆΠΈΡ‚ подписки ΠΊΠ»ΠΈΠ΅Π½Ρ‚ΠΎΠ² ΠΏΠΎ room/userId. Worker Ρ‡ΠΈΡ‚Π°Π΅Ρ‚ change stream, ΠΌΠ°Ρ€ΡˆΡ€ΡƒΡ‚ΠΈΠ·ΠΈΡ€ΡƒΠ΅Ρ‚ событиС Π² Π½ΡƒΠΆΠ½Ρ‹Π΅ room. Π’Π°ΠΆΠ½ΠΎ: fan-out Π½Π° тысячи ΠΊΠ»ΠΈΠ΅Π½Ρ‚ΠΎΠ² Π½Π΅ Π΄Π΅Ρ€ΠΆΠΈΡ‚Π΅ Π² ΠΎΠ΄Π½ΠΎΠΌ процСссС Π±Π΅Π· horizontal scale gateway.

Idempotency ΠΈ ordering

Consumer ΠΌΠΎΠΆΠ΅Ρ‚ ΠΏΠΎΠ»ΡƒΡ‡ΠΈΡ‚ΡŒ duplicate ΠΏΡ€ΠΈ reconnect. Π₯Ρ€Π°Π½ΠΈΡ‚Π΅ processed event id (hash clusterTime + documentKey + operationType) Π² Redis ΠΈΠ»ΠΈ локальной Ρ‚Π°Π±Π»ΠΈΡ†Π΅ с TTL. БизнСс-ΠΎΠ±Ρ€Π°Π±ΠΎΡ‚Ρ‡ΠΈΠΊ Π΄ΠΎΠ»ΠΆΠ΅Π½ Π±Ρ‹Ρ‚ΡŒ ΠΈΠ΄Π΅ΠΌΠΏΠΎΡ‚Π΅Π½Ρ‚Π½Ρ‹ΠΌ: Β«ΠΎΠ±Π½ΠΎΠ²ΠΈΡ‚ΡŒ search indexΒ» бСзопасно ΠΏΠΎΠ²Ρ‚ΠΎΡ€ΠΈΡ‚ΡŒ, Β«ΡΠΏΠΈΡΠ°Ρ‚ΡŒ баланс» β€” Π½Π΅Ρ‚ Π±Π΅Π· dedup.

Backpressure

ΠŸΡ€ΠΈ burst записСй Π² MongoDB consumer Π½Π΅ успСваСт β€” lag растёт, memory pressure. РСшСния: batch ΠΎΠ±Ρ€Π°Π±ΠΎΡ‚ΠΊΠ°, ΠΎΡ‚Π΄Π΅Π»ΡŒΠ½Π°Ρ ΠΎΡ‡Π΅Ρ€Π΅Π΄ΡŒ с consumer group, ΠΌΠ°ΡΡˆΡ‚Π°Π±ΠΈΡ€ΠΎΠ²Π°Π½ΠΈΠ΅ workers ΠΏΠΎ partition key (shard _id range).

ВрСбования, ограничСния ΠΈ вСрсии

  • Replica set обязатСлСн; standalone Π½Π΅ ΠΏΠΎΠ΄Π΄Π΅Ρ€ΠΆΠΈΠ²Π°Π΅Ρ‚ change streams.
  • Sharded cluster: stream Π½Π° sharded collection ΡΠΎΠ±Π»ΡŽΠ΄Π°Π΅Ρ‚ порядок Ρ‚ΠΎΠ»ΡŒΠΊΠΎ для Π΄ΠΎΠΊΡƒΠΌΠ΅Π½Ρ‚ΠΎΠ² с ΠΎΠ΄Π½ΠΈΠΌ shard key value; Π³Π»ΠΎΠ±Π°Π»ΡŒΠ½Ρ‹ΠΉ total order Π½Π΅ Π³Π°Ρ€Π°Π½Ρ‚ΠΈΡ€ΠΎΠ²Π°Π½.
  • Transaction events Π²ΠΈΠ΄Π½Ρ‹ ΠΊΠ°ΠΊ ΠΎΡ‚Π΄Π΅Π»ΡŒΠ½Ρ‹Π΅ ΠΎΠΏΠ΅Ρ€Π°Ρ†ΠΈΠΈ ΠΈΠ»ΠΈ Ρ‡Π΅Ρ€Π΅Π· full transaction Π² зависимости ΠΎΡ‚ вСрсии.
  • Pre/post images (change stream pre/post images) Π΄Π°ΡŽΡ‚ Π΄ΠΎ/послС для update β€” ΠΏΠΎΠ»Π΅Π·Π½ΠΎ для audit, Ρ‚Ρ€Π΅Π±ΡƒΠ΅Ρ‚ Π²ΠΊΠ»ΡŽΡ‡Π΅Π½ΠΈΡ Π½Π° ΠΊΠΎΠ»Π»Π΅ΠΊΡ†ΠΈΠΈ.
  • DDL (drop collection) Π³Π΅Π½Π΅Ρ€ΠΈΡ€ΡƒΠ΅Ρ‚ invalidate β€” consumer Π΄ΠΎΠ»ΠΆΠ΅Π½ ΠΏΠ΅Ρ€Π΅ΡΠΎΠ·Π΄Π°Ρ‚ΡŒ stream.

На Atlas доступны Π΄ΠΎΠΏΠΎΠ»Π½ΠΈΡ‚Π΅Π»ΡŒΠ½Ρ‹Π΅ ΠΈΠ½Ρ‚Π΅Π³Ρ€Π°Ρ†ΠΈΠΈ; self-hosted Ρ‚Ρ€Π΅Π±ΡƒΠ΅Ρ‚ ΠΌΠΎΠ½ΠΈΡ‚ΠΎΡ€ΠΈΠ½Π³Π° oplog window: Ссли consumer отстаёт дольшС, Ρ‡Π΅ΠΌ хранится oplog, resume Π½Π΅Π²ΠΎΠ·ΠΌΠΎΠΆΠ΅Π½ Π±Π΅Π· full sync.

Π‘Π΅Π·ΠΎΠΏΠ°ΡΠ½ΠΎΡΡ‚ΡŒ ΠΈ эксплуатация

Change stream cursor ΠΈΡΠΏΠΎΠ»ΡŒΠ·ΡƒΠ΅Ρ‚ ΠΏΡ€Π°Π²Π° чтСния Π½Π° Π±Π°Π·Ρƒ. Π’Ρ‹Π΄Π΅Π»ΠΈΡ‚Π΅ role readChangeStream + read Π½Π° Π½ΡƒΠΆΠ½Ρ‹Π΅ ΠΊΠΎΠ»Π»Π΅ΠΊΡ†ΠΈΠΈ для service account worker, Π½Π΅ admin. TLS ΠΌΠ΅ΠΆΠ΄Ρƒ worker ΠΈ MongoDB обязатСлСн Π² prod.

ΠœΠ΅Ρ‚Ρ€ΠΈΠΊΠΈ для Π°Π»Π΅Ρ€Ρ‚ΠΎΠ²:

  • lag ΠΌΠ΅ΠΆΠ΄Ρƒ clusterTime события ΠΈ Π²Ρ€Π΅ΠΌΠ΅Π½Π΅ΠΌ ΠΎΠ±Ρ€Π°Π±ΠΎΡ‚ΠΊΠΈ;
  • rate ошибок resume (ExpiredChangeStreamEvent);
  • Ρ€Π°Π·ΠΌΠ΅Ρ€ ΠΎΡ‡Π΅Ρ€Π΅Π΄ΠΈ downstream;
  • restarts consumer pod Π±Π΅Π· сохранённого token.

ВСстируйтС failover: primary step-down Π½Π΅ Π΄ΠΎΠ»ΠΆΠ΅Π½ ΡƒΠ±ΠΈΠ²Π°Ρ‚ΡŒ stream навсСгда β€” Π΄Ρ€Π°ΠΉΠ²Π΅Ρ€ ΠΏΠ΅Ρ€Π΅ΠΏΠΎΠ΄ΠΊΠ»ΡŽΡ‡Π°Π΅Ρ‚ΡΡ ΠΊ Π½ΠΎΠ²ΠΎΠΌΡƒ primary с resume token.

ΠŸΡ€Π°ΠΊΡ‚ΠΈΡ‡Π΅ΡΠΊΠΈΠΉ сцСнарий: live-Π΄Π°ΡˆΠ±ΠΎΡ€Π΄ Π·Π°ΠΊΠ°Π·ΠΎΠ²

ΠšΠΎΠ»Π»Π΅ΠΊΡ†ΠΈΡ orders с полями status, warehouseId, updatedAt. Pipeline: $match status in [processing, shipped]. Worker ΠΏΠΈΡˆΠ΅Ρ‚ Π² Redis pub/sub channel warehouse:{id}. Frontend подписан Π½Π° WebSocket ΠΈ обновляСт Ρ‚Π°Π±Π»ΠΈΡ†Ρƒ Π±Π΅Π· refresh. ΠŸΡ€ΠΈ update status Π½Π° delivered событиС ΡƒΡ…ΠΎΠ΄ΠΈΡ‚ Ρ‚ΠΎΠ»ΡŒΠΊΠΎ подписчикам Π½ΡƒΠΆΠ½ΠΎΠ³ΠΎ склада.

Π”ΠΎΠΏΠΎΠ»Π½ΠΈΡ‚Π΅Π»ΡŒΠ½ΠΎ: Ρ‚ΠΎΡ‚ ΠΆΠ΅ stream ΠΊΠΎΡ€ΠΌΠΈΡ‚ Elasticsearch для полнотСкстового поиска; idempotency ΠΏΠΎ orderId + updatedAt version field ΠΏΡ€Π΅Π΄ΠΎΡ‚Π²Ρ€Π°Ρ‰Π°Π΅Ρ‚ stale overwrite.

Change Streams vs Π°Π»ΡŒΡ‚Π΅Ρ€Π½Π°Ρ‚ΠΈΠ²Ρ‹

Polling β€” просто, Π½ΠΎ latency ΠΈ load Π½Π° DB. Triggers Π² ΠΏΡ€ΠΈΠ»ΠΎΠΆΠ΅Π½ΠΈΠΈ β€” Π»Π΅Π³ΠΊΠΎ ΠΏΡ€ΠΎΠΏΡƒΡΡ‚ΠΈΡ‚ΡŒ ΠΏΡƒΡ‚ΡŒ записи. Debezium CDC β€” ΠΌΠΎΡ‰Π½Π΅Π΅ для multi-DB, Π½ΠΎ тяТСлСС Π² ops. Change Streams β€” sweet spot Π²Π½ΡƒΡ‚Ρ€ΠΈ MongoDB-only landscape с ΠΌΠΈΠ½ΠΈΠΌΠ°Π»ΡŒΠ½Ρ‹ΠΌ glue.

Если real-time слой ΠΊΡ€ΠΈΡ‚ΠΈΡ‡Π΅Π½ для ΠΏΡ€ΠΎΠ΄ΡƒΠΊΡ‚Π°, Π·Π°ΠΊΠ»Π°Π΄Ρ‹Π²Π°ΠΉΡ‚Π΅ Π΅Π³ΠΎ Π² Π°Ρ€Ρ…ΠΈΡ‚Π΅ΠΊΡ‚ΡƒΡ€Ρƒ с ΠΏΠ΅Ρ€Π²ΠΎΠ³ΠΎ Ρ€Π΅Π»ΠΈΠ·Π°: resume tokens, idempotency, ΠΌΠΎΠ½ΠΈΡ‚ΠΎΡ€ΠΈΠ½Π³ oplog lag β€” Π½Π΅ Β«Π΄ΠΎΠ±Π°Π²ΠΈΠΌ ΠΏΠΎΡ‚ΠΎΠΌΒ».

Π Π°Π·Π±ΠΎΡ€ ΠΈΠ½Ρ‚Π΅Π³Ρ€Π°Ρ†ΠΈΠΈ Change Streams Π² ваш стСк β€” Π½Π° ΠΊΠΎΠ½ΡΡƒΠ»ΡŒΡ‚Π°Ρ†ΠΈΠΈ ΠΏΠΎ real-time MongoDB.

НадёТный consumer Change Streams

Change Streams ΡƒΠ΄ΠΎΠ±Π½Ρ‹ для fan-out событий, Π½ΠΎ Π² ΠΏΡ€ΠΎΠ΄Π΅ Π²Π°ΠΆΠ½Ρ‹ resume token, ΠΎΠ±Ρ€Π°Π±ΠΎΡ‚ΠΊΠ° ошибок ΠΈ backpressure. Π₯Ρ€Π°Π½ΠΈΡ‚Π΅ token Π² Π½Π°Π΄Ρ‘ΠΆΠ½ΠΎΠΌ store; послС рСстарта consumer Π΄ΠΎΠ»ΠΆΠ΅Π½ ΠΏΡ€ΠΎΠ΄ΠΎΠ»ΠΆΠΈΡ‚ΡŒ с Ρ‚ΠΎΠ³ΠΎ ΠΆΠ΅ мСста, Π° Π½Π΅ «с сСйчас».

НС Ρ‚Π°Ρ‰ΠΈΡ‚Π΅ Ρ‚ΡΠΆΡ‘Π»ΡƒΡŽ бизнСс-Π»ΠΎΠ³ΠΈΠΊΡƒ прямо Π² колбэк стрима: ΠΊΠ»Π°Π΄ΠΈΡ‚Π΅ события Π² ΠΎΡ‡Π΅Ρ€Π΅Π΄ΡŒ ΠΈ ΠΎΠ±Ρ€Π°Π±Π°Ρ‚Ρ‹Π²Π°ΠΉΡ‚Π΅ Π²ΠΎΡ€ΠΊΠ΅Ρ€Π°ΠΌΠΈ. Π˜Π½Π°Ρ‡Π΅ ΠΎΠ΄ΠΈΠ½ ΠΌΠ΅Π΄Π»Π΅Π½Π½Ρ‹ΠΉ handler остановит Ρ‡Ρ‚Π΅Π½ΠΈΠ΅.

  • Π€ΠΈΠ»ΡŒΡ‚Ρ€Ρ‹ pipeline Π½Π° Π½ΡƒΠΆΠ½Ρ‹Π΅ operationType ΠΈ поля
  • Π˜Π΄Π΅ΠΌΠΏΠΎΡ‚Π΅Π½Ρ‚Π½ΠΎΡΡ‚ΡŒ ΠΏΠΎ ΠΈΠ΄Π΅Π½Ρ‚ΠΈΡ„ΠΈΠΊΠ°Ρ‚ΠΎΡ€Ρƒ события
  • АлСрты Π½Π° Π»Π°Π³ consumer ΠΈ ΠΏΠΎΡ‚Π΅Ρ€ΠΈ соСдинСния
  • ВСст failover replica set: стрим Π΄ΠΎΠ»ΠΆΠ΅Π½ ΠΏΠ΅Ρ€Π΅ΠΆΠΈΡ‚ΡŒ primary stepDown

Π“Ρ€Π°Π½ΠΈΡ†Ρ‹ примСнимости

Change Streams β€” Π½Π΅ Π·Π°ΠΌΠ΅Π½Π° ΠΏΠΎΠ»Π½ΠΎΡ†Π΅Π½Π½ΠΎΠΉ event-driven Π°Ρ€Ρ…ΠΈΡ‚Π΅ΠΊΡ‚ΡƒΡ€Π΅. Для ΠΎΡ‡Π΅Π½ΡŒ высоких write TPS ΠΈ слоТного Ρ„Π°Π½-Π°ΡƒΡ‚Π° часто Π²Ρ‹Π³ΠΎΠ΄Π½Π΅Π΅ явный outbox. Π˜ΡΠΏΠΎΠ»ΡŒΠ·ΡƒΠΉΡ‚Π΅ стримы Ρ‚Π°ΠΌ, Π³Π΄Π΅ MongoDB ΡƒΠΆΠ΅ source of truth ΠΈ Π½ΡƒΠΆΠ΅Π½ ΠΎΡ‚Π½ΠΎΡΠΈΡ‚Π΅Π»ΡŒΠ½ΠΎ простой real-time слой.

Π—Π°ΠΏΠΈΡΠ°Ρ‚ΡŒΡΡ Π½Π° Π±Π΅ΡΠΏΠ»Π°Ρ‚Π½ΡƒΡŽ ΠΊΠΎΠ½ΡΡƒΠ»ΡŒΡ‚Π°Ρ†ΠΈΡŽ.