MongoDB: تدفقات التغييرات: مراقبة التغييرات في الوقت الفعلي
تُعد «تدفقات التغييرات» (Change Streams) الأساس الذي تستند إليه إشعارات التغييرات في الوقت الفعلي في MongoDB — المعروفة باسم CDC (التقاط بيانات التغيير) — والتطبيقات التي تعمل استجابةً للأحداث.
1. ما ستتعلمه
- مفهوم تيارات التغيير (استنادًا إلى oplog)
- db.collection.watch() الاستماع
- تصفية خط الأنابيب ($match)
- fullDocument: 'updateLookup'
- التكامل مع Mongoose
- السيناريوهات النموذجية (الإشعارات في الوقت الفعلي، مزامنة البيانات)
graph LR
A[Applications 1<br/>Insert Order] -->|oplog| DB[(MongoDB<br/>oplog)]
B[Applications 2<br/>Update Order] -->|oplog| DB
C[Applications 3<br/>Delete Order] -->|oplog| DB
DB -->|Change Stream| CS[db.collection.watch]
CS -->|operationType=insert| N1[Email Notification<br/>WebSocket Push]
CS -->|operationType=update| N2[Elasticsearch<br/>Synchronize]
CS -->|operationType=delete| N3[Redis Cache<br/>Invalidation]
style CS fill:#d4edda
style N1 fill:#cce5ff
2. أساسيات «تغيير التدفقات»
نظرة عامة على المفهوم: «Change Streams» هي واجهة برمجة تطبيقات (API) لمراقبة التغييرات في الوقت الفعلي، تم تقديمها في الإصدار 3.6 من MongoDB والإصدارات الأحدث، ويتم تنفيذها باستخدام سجل عمليات (oplog) مجموعة النسخ المتماثلة. وهي تتيح للتطبيقات اكتشاف التغييرات في البيانات في الوقت الفعلي دون الحاجة إلى الاستقصاء، مما يجعلها الأساس الذي تستند إليه تقنية «CDC» (التقاط بيانات التغيير) والبنى المعمارية التي تعتمد على الأحداث.
مبادئ OpLog: يُعد Oplog (سجل العمليات) مجموعة خاصة محدودة السعة تسجل جميع عمليات الكتابة في مجموعة النسخ المتماثلة. يقوم المكون الأساسي بإضافة كل عملية كتابة (إدراج/تحديث/حذف) إلى Oplog، بينما يقوم المكون الثانوي بنسخ البيانات عن طريق متابعة Oplog. وتعد تيارات التغيير (Change Streams) في الأساس تغليفًا منظمًا لـ Oplog — فهي تحول إدخالات Oplog الأولية إلى أحداث تغيير يمكن قراءتها بواسطة البشر، وتوفر ميزات متقدمة مثل التصفية ووظيفة الاستئناف من نقطة التوقف.
graph TB
subgraph "OpLog How It Works"
W1[Write Operation] --> OP[(oplog<br/>Fixed Set)]
OP --> S1[Secondary 1<br/>tailing oplog]
OP --> S2[Secondary 2<br/>tailing oplog]
OP --> CS[Change Stream<br/>Structured Events]
end
subgraph "Change Stream vs Polling"
P1[Polling: Queries per second] --> C1[High latency<br/>High Load<br/>Duplicate Query]
CS --> C2[Real-time Notifications<br/>Low load<br/>No duplicates]
end
style CS fill:#d4edda
style C2 fill:#d4edda
style C1 fill:#f8d7da
تدفقات التغيير مقابل الاستقصاء: مقارنة:
| البعد | الاستطلاعات | تغييرات |
|---|---|---|
| التأخير | 1–60 ثانية (حسب الفاصل الزمني) | < 100 مللي ثانية (إرسال فوري) |
| حمل قاعدة البيانات | مرتفع (مسح الفهرس في كل استعلام) | منخفض (تلقي أحداث سجل العمليات بشكل سلبي) |
| سلامة الأحداث | احتمال فقدان الأحداث (تغييرات متعددة خلال فترة زمنية معينة) | لا يوجد فقدان (سجل الأحداث مرتب) |
| تنزيل ملف الاستئناف | يجب تنفيذه يدويًّا | رمز الاستئناف المدمج |
| استخدام الموارد | استخدام الاتصال المستمر + وحدة المعالجة المركزية | اتصالات طويلة الأمد، استخدام منخفض لوحدة المعالجة المركزية |
تغيير بنية الحدث:
// === Monitor Set Changes ===
const changeStream = db.products.watch();
changeStream.on('change', (change) => {
console.log('Change detected:', change);
// {
// _id: { _data: '...' }, // resumeToken(For resuming downloads)
// operationType: 'insert', // Operation Type
// fullDocument: { _id: ..., sku: ..., title: ..., ... }, // Complete Documentation
// ns: { db: 'shopdb', coll: 'products' }, // Namespace
// documentKey: { _id: ObjectId('...') } // Document Primary Key
// }
});
| نوع_العملية | المعنى | الوثيقة_الكاملة |
|---|---|---|
insert |
إدراج مستند جديد | ✅ مستند كامل |
update |
تحديث المستند | ❌ تغييرات في الحقول فقط (يتطلب استخدام updateLookup) |
replace |
استبدال المستند | ✅ المستند الكامل |
delete |
تم حذف المستند | ❌ لا شيء (تم حذف المستند) |
drop |
حذف المجموعة | ❌ لا شيء |
rename |
مجموعة إعادة التسمية | ❌ لا شيء |
invalidate |
فقدان التأثير (مثل حذف المجموعة) | ❌ لا شيء |
تحليل النقاط الرئيسية:
- يجب أن تعمل «تدفقات التغيير» على مجموعة النسخ المتماثلة أو الكتلة المقسمة (حسب سجل العمليات).
updateبشكل افتراضي، لا يُرجع الحدث المستند بالكامل؛ يجب عليك تعيينfullDocument: 'updateLookup'- الحقل
_idهو رمز استئناف التنزيل (resumeToken)، الذي يُستخدم لاستئناف التنزيلات — احفظه، وبعد إعادة التشغيل، تابع التنزيل من حيث توقفت.
3. ترشيح خطوط الأنابيب
وصف المفهوم: يدعم «Change Streams» تصفية مسار التجميع — من خلال تطبيق مراحل مثل $match و$project على دفق أحداث oplog، بحيث لا يتم تمرير سوى الأحداث ذات الأهمية إلى التطبيق. تتم عملية التصفية من جانب الخادم، مما يقلل من حركة مرور الشبكة والأعباء الإضافية للمعالجة على مستوى طبقة التطبيق.
كيفية العمل: تقبل watch() مصفوفة من مسارات التجميع كمعلمة؛ ويتم تنفيذ مراحل المسار على جانب خادم MongoDB. تمر الأحداث أولاً عبر خط الأنابيب للتصفية، ولا يتم إرسال سوى تلك التي تجتاز التصفية إلى العميل. وهو يدعم مراحل مثل $match و$project و$addFields و$replaceRoot.
سياسة التصفية:
| هدف التصفية | مرحلة مسار المعالجة | مثال |
|---|---|---|
| نوع العملية | $match: {operationType} |
الاستماع إلى عمليات الإدراج فقط |
| حقل المستند | $match: {'fullDocument.field'} |
الفئات المحددة فقط |
| الحقل المراد تغييره | $match: {'updateDescription.updatedFields'} |
تغيير السعر فقط |
| معايير التجميع | $match: {$or: [...]} |
الإضافة أو تغيير السعر |
// === Filter by Specific Actions ===
const changeStream = db.products.watch([
{ $match: { operationType: 'insert' } }
]);
// === Filter by Specific Fields ===
const changeStream = db.products.watch([
{ $match: { 'fullDocument.category': 'Electronics' } }
]);
// === Filter by Price Changes ===
const changeStream = db.products.watch([
{
$match: {
$or: [
{ operationType: 'insert' },
{ operationType: 'update', 'updateDescription.updatedFields.price': { $exists: true } }
]
}
}
]);
تحليل النقاط الرئيسية:
- تتم عملية التصفية من جانب الخادم، مما يقلل من حركة مرور الشبكة غير الضرورية.
- يستخدم مسار الحقل الخاص بـ
$matchبنية أحداث oplog (مثلfullDocument.category)، وليس حقول المستند الأصلي. - المراحل التي تتطلب حالة عامة، مثل
$groupو$limit، غير مدعومة.
4. تكوين fullDocument
وصف المفهوم: يتحكم الخيار fullDocument في ما إذا كانت أحداث «Change Stream» تتضمن المستند الكامل أم لا. بشكل افتراضي، لا تعرض أحداث update سوى الحقول التي طرأت عليها تغييرات (updateDescription) ولا تعرض المستند الكامل. وعند تعيين fullDocument: 'updateLookup'، يقوم MongoDB بإجراء استعلام إضافي لاسترداد المحتوى الكامل للمستند الحالي.
كيف يعمل:
whenAvailable(الافتراضي): تشمل أحداث الإدراج/الاستبدال المستند بأكمله؛ بينما تشمل أحداث التحديث الحقول التي طرأت عليها تغييرات فقط؛ أما أحداث الحذف فلا تشمل أي مستندupdateLookup: تشمل جميع الأحداث الوثيقة الكاملة (وستقوم أحداث التحديث بالاستعلام عن أحدث نسخة من الوثيقة بالإضافة إلى ذلك)، ولكن هناك عبء إضافي على عملية القراءة وتأخير طفيف.required: يجب إرجاع المستند كاملاً؛ وإلا فسيتم الإبلاغ عن خطأ.
المزايا والعيوب في updateLookup:
| البعد | وقت التوفر | تحديث البحث |
|---|---|---|
| الوثائق الكاملة | الإدراج/الاستبدال فقط | جميع أنواع العمليات |
| الأداء | معيار الأداء | العبء الإضافي للاستعلامات (+10–20٪) |
| حداثة البيانات | وقت التغيير | وقت الاستعلام (قد يتأخر قليلاً) |
| حالات الاستخدام | الحاجة إلى معرفة الحقول التي طرأت عليها تغييرات فقط | الحاجة إلى المستند الكامل من أجل المعالجة اللاحقة |
sequenceDiagram
participant App as Applications
participant DB as MongoDB
participant Doc as Document
App->>DB: watch([], {fullDocument: 'updateLookup'})
DB->>DB: oplog Generate update Event
Note over DB: update By default, events contain only the changed fields.
DB->>Doc: View the current complete document
Doc-->>DB: Back to Latest Documents
DB-->>App: Push Events + fullDocument
Note over App: The incident includes the complete documentation
// === update Return to the full document ===
const changeStream = db.products.watch([], {
fullDocument: 'updateLookup'
});
changeStream.on('change', (change) => {
if (change.operationType === 'update') {
console.log('Updated doc:', change.fullDocument);
// Complete Documentation(After the default value changes)
}
});
// === Return only the fields that have changed ===
const changeStream = db.products.watch([], {
fullDocument: 'whenAvailable' // Default
});
تحليل النقاط الرئيسية:
updateLookupيُجري استعلامًا إضافيًا، مما يزيد من الحمل على قاعدة البيانات في الحالات التي تتكرر فيها عمليات التحديث.- تُرجع
updateLookupأحدث مستند حتى وقت الاستعلام، والذي قد يختلف قليلاً عن المستند الموجود وقت إجراء التغيير (بسبب تعديلات متزامنة أخرى). - لا يُرجع الحدث
deleteمستندًا حتى لو تم تعيينupdateLookup(لم يعد المستند موجودًا)
5. تكامل Mongoose
شرح المفهوم: يدعم Mongoose 6+ ميزة «تدفقات التغييرات» (Change Streams) بشكل أصلي، والتي يتم إرجاعها عبر Model.watch(). ويتوافق هذا تمامًا مع واجهة برمجة التطبيقات (API) الخاصة ببرنامج التشغيل الأصلي لـ MongoDB، لكن استخدامها مباشرةً في طبقة النموذج يتماشى بشكل أفضل مع قواعد التطوير المتبعة في Mongoose.
البنية القائمة على الأحداث: تُعد «تدفقات التغييرات» مكونًا أساسيًّا في البنية القائمة على الأحداث — حيث تُعد التغييرات التي تطرأ على قاعدة البيانات مصدرًا للأحداث، مما يؤدي إلى آثار جانبية لاحقة مثل تحديثات ذاكرة التخزين المؤقت، ومزامنة عمليات البحث، وإرسال الإشعارات.
graph TB
subgraph "Source of the Incident"
DB[(MongoDB<br/>oplog)]
end
subgraph "Change Stream Bus"
CS[Model.watch<br/>Change Flow]
end
subgraph "Event Consumer"
N1[Email Notification]
N2[Elasticsearch<br/>Search Synchronization]
N3[Redis<br/>Cache Expiration]
N4[WebSocket<br/>Real-time Notifications]
N5[Audit Log<br/>Compliance Record]
end
DB --> CS
CS --> N1
CS --> N2
CS --> N3
CS --> N4
CS --> N5
style CS fill:#d4edda
// === mongoose Change Streams(mongoose 6+)===
const Product = mongoose.model('Product', productSchema);
// Monitoring Product Set Changes
const changeStream = Product.watch();
changeStream.on('change', (change) => {
console.log(`${change.operationType}:`, change.fullDocument);
});
// Filter
const filteredStream = Product.watch([
{ $match: { operationType: { $in: ['insert', 'update'] } } }
]);
تحليل النقاط الرئيسية:
- يُرجع
watch()من Mongoose تيار التغييرات (Change Stream) الأصلي في MongoDB، وتكون واجهة برمجة التطبيقات (API) مطابقة تمامًا. - يجب الحفاظ على اتصال قاعدة البيانات أثناء عملية المراقبة؛ فإذا انقطع الاتصال، يتم إيقاف تشغيل Change Stream تلقائيًّا.
- في بيئة الإنتاج، يُوصى بتغليف «مدير تيار التغييرات» (Change Stream Manager): إعادة الاتصال التلقائي + استمرارية رمز الاستئناف (resumeToken)
6. سيناريوهات واقعية
نظرة عامة على المفهوم: توجد ثلاث حالات استخدام تقليدية لـ «تدفقات التغيير» (Change Streams) في بيئات الإنتاج: الإشعارات في الوقت الفعلي، ومزامنة البيانات (CDC)، وإبطال صلاحية ذاكرة التخزين المؤقت. وتُعد كل حالة من حالات الاستخدام هذه نموذجًا نموذجيًّا لتنفيذ بنية قائمة على الأحداث.
(1) نظام الإخطار في الوقت الفعلي
السيناريو: تحتاج منصة ShopHub للتجارة الإلكترونية إلى إرسال إشعارات عبر البريد الإلكتروني في الوقت الفعلي وإشعارات دفع عبر WebSocket إلى النظام الخلفي للتاجر عند إنشاء طلب جديد.
// === Monitor New Orders,Send a notification ===
const OrderStream = db.orders.watch([
{ $match: { operationType: 'insert' } }
]);
OrderStream.on('change', async (change) => {
const order = change.fullDocument;
// Send an Email
await sendEmail(order.userId, 'Order confirmation', `Order ${order._id} received`);
// Push Notifications
await pushNotification(order.userId, {
title: 'New Order',
body: `Order total: $${order.total}`
});
// WebSocket Real-time Notifications
io.emit('new_order', order);
});
(2) مزامنة البيانات (CDC)
السيناريو: تحتاج ShopHub إلى مزامنة بيانات المنتجات الموجودة في MongoDB في الوقت الفعلي مع Elasticsearch لدعم البحث عن النص الكامل. تتيح تقنية Change Streams إنشاء مسار نقل بيانات التغييرات (CDC) بدون أي تأخير.
الهندسة المعمارية لمركز السيطرة على الأمراض (CDC):
graph LR
A[MongoDB<br/>Product Data] -->|Change Stream| B[CDC Worker<br/>Node.jsProcess]
B -->|insert/update| C[Elasticsearch<br/>Search Index]
B -->|delete| C
B -->|Change Log| D[(Redis<br/>resumeToken)]
D -->|Restart and Recover| B
style B fill:#d4edda
// === Monitoring MongoDB Change,Sync to Elasticsearch ===
const ProductStream = db.products.watch();
ProductStream.on('change', async (change) => {
switch (change.operationType) {
case 'insert':
case 'update':
case 'replace':
await elasticsearch.index({
index: 'products',
id: change.documentKey._id.toString(),
body: change.fullDocument
});
break;
case 'delete':
await elasticsearch.delete({
index: 'products',
id: change.documentKey._id.toString()
});
break;
}
});
(3) إبطال صلاحية ذاكرة التخزين المؤقت
السيناريو: يستخدم ShopHub Redis لتخزين صفحات تفاصيل المنتجات مؤقتًا. وعندما تتغير بيانات المنتج، يتم مسح ذاكرة التخزين المؤقت تلقائيًا لمنع وجود بيانات قديمة.
// === Monitor Product Changes,Invalidate Redis cache ===
const ProductStream = db.products.watch();
ProductStream.on('change', async (change) => {
const productId = change.documentKey._id.toString();
await redis.del(`product:${productId}`);
console.log(`Cache cleared for ${productId}`);
});
▶ المثال 1: استئناف تدفق التغييرات بعد الانقطاع
// Alice's TechCorp System: Change Stream Resume, Events Are Not Lost After a Restart
async function startResumableStream() {
// 1. Retrieve last saved resumeToken from Redis
let resumeToken = await redis.get('product_stream_token');
let options = { fullDocument: 'updateLookup' };
if (resumeToken) {
options.resumeAfter = JSON.parse(resumeToken);
console.log('Resuming from saved token');
}
// 2. Start Listening
const changeStream = db.products.watch([], options);
changeStream.on('change', async (change) => {
// Processing Changes...
console.log(`${change.operationType}: ${change.documentKey._id}`);
// 3. Save after each event resumeToken
await redis.set('product_stream_token', JSON.stringify(change._id));
});
changeStream.on('error', async (err) => {
console.error('Stream error:', err.message);
// 4. Delayed Reconnection After an Error
setTimeout(startResumableStream, 5000);
});
}
startResumableStream();
الإخراج:
TEXT 📖 للعرض فقطResuming from saved token insert: 647f1f77bcf86cd799439011 update: 647f1f77bcf86cd799439012
7. قيود تيار التغييرات
ملاحظة مفاهيمية: تعتمد «تدفقات التغيير» (Change Streams) على سجل التغييرات (oplog) ومجموعات النسخ المتماثلة، ولها قيود واضحة. ويُعد فهم هذه القيود أمرًا ضروريًا لتصميم أنظمة موثوقة تعمل في الوقت الفعلي.
شرح مفصل للقيود:
| القيد | الوصف | السبب | استراتيجية التخفيف |
|---|---|---|---|
| يلزم وجود مجموعة نسخ متماثلة | غير مدعوم في الوضع المستقل | يعتمد على سجل العمليات (oplog) | استخدم مجموعة نسخ متماثلة ذات عقدة واحدة في بيئة التطوير |
| الحد الأقصى لحجم سجل العمليات (oplog) | القيمة الافتراضية: 5% من مساحة القرص | سجل العمليات (oplog) هو مجموعة محددة السعة | قم بزيادة قيمة oplogSize أو راقب النافذة |
| لا يمكن أن يمتد عبر المجموعات | داخل مجموعة واحدة | لا يمتد سجل العمليات (Oplog) عبر المجموعات | استخدم Kafka لربط مجموعات متعددة |
| $where غير مدعوم | بعض العوامل | لا يسجل Oplog تفاصيل الاستعلام | التصفية باستخدام حقل fullDocument |
| ترتيب الأحداث | يتم الترتيب داخل المجموعة الواحدة؛ ولا يوجد ترتيب شامل عبر المجموعات | يتم تقسيم سجل العمليات (Oplog) حسب المجموعة | يتم إجراء الفرز عبر المجموعات على مستوى طبقة التطبيق |
| استخدام الذاكرة | الذاكرة المستخدمة لكل اتصال مراقبة | يحافظ الخادم على حالة المؤشر | الحد الأقصى لعدد اتصالات المراقبة |
مراقبة نافذة سجل التغييرات (oplog): يتميز سجل التغييرات (oplog) بحجم ثابت، ويتم الكتابة فوق الإدخالات الأقدم. إذا كان معدل استهلاك البيانات بواسطة «Change Stream» أبطأ من معدل الكتابة فوق سجل التغييرات (oplog)، فسوف تنتهي صلاحية «resumeToken»، ولن تعمل وظيفة «الاستئناف من نقطة التوقف».
graph LR
A[oplog Write Speed<br/>1000 ops/s] --> B[oplog Capacity<br/>5% Disk ≈ 50GB]
B --> C[oplog Window<br/>about 72 hours]
C --> D{Consumption Rate?}
D -->|Keep up| E[✅ Normal]
D -->|Behind > 72h| F[❌ resumeToken Failure<br/>A full resync is required]
style E fill:#d4edda
style F fill:#f8d7da
| عنصر المراقبة | الأمر | عتبة التنبيه |
|---|---|---|
| نافذة سجل العمليات | rs.printReplicationInfo() |
أقل من 24 ساعة |
| استخدام سجل العمليات (oplog) | db.oplog.rs.stats() |
> 80% |
| عدد اتصالات مسارات التغيير | db.currentOp() |
> 100 |
| تأخير الحدث | مراقبة طبقة التطبيق | > 10 ثوانٍ |
▶ مثال: دليل عملي لتعديل إشعارات الطلبات في الوقت الفعلي في «Streams»
// Scene: Monitor New Orders, Send emails in real time + WebSocket Push + Sync to Elasticsearch
// 1. Start Listening (in Node.js)
const { MongoClient } = require('mongodb');
async function startOrderListener() {
const client = new MongoClient('mongodb://localhost:27017/?replicaSet=rs0');
await client.connect();
const orders = client.db('shopdb').collection('orders');
// Monitor all changes to the order collection
const changeStream = orders.watch([
{
$match: {
operationType: 'insert', // Listen-only insertion
'fullDocument.status': 'paid' // Process only paid orders
}
}
]);
changeStream.on('change', async (change) => {
const order = change.fullDocument;
console.log(`New Orders: ${order._id}, Amount: $${order.total}`);
// 1. Send an email notification
await sendEmail(order.userId, {
subject: 'Order Confirmation',
body: `Your Order ${order._id} Submitted,Total Amount $${order.total}`
});
// 2. WebSocket Real-time push notifications to the merchant's backend
io.to('merchant-dashboard').emit('new_order', {
orderId: order._id,
total: order.total,
items: order.items,
timestamp: order.createdAt
});
// 3. Sync to Elasticsearch For Search
await elasticsearch.index({
index: 'orders',
id: order._id.toString(),
body: order
});
});
// 4. Monitor Product Changes,Synchronized Update Elasticsearch
const productStream = client.db('shopdb').collection('products').watch();
productStream.on('change', async (change) => {
if (change.operationType === 'delete') {
await elasticsearch.delete({
index: 'products',
id: change.documentKey._id.toString()
});
} else if (change.fullDocument) {
await elasticsearch.index({
index: 'products',
id: change.documentKey._id.toString(),
body: change.fullDocument
});
}
});
console.log('Change Streams listening...');
}
// 5. Resume Download(Restore progress tracking after the app restarts)
async function resumeAfterRestart() {
const resumeToken = await redis.get('change_stream_resume_token');
const changeStream = orders.watch([], {
resumeAfter: JSON.parse(resumeToken),
fullDocument: 'updateLookup'
});
changeStream.on('change', async (change) => {
// Processing Changes...
// Save resume token
await redis.set('change_stream_resume_token', JSON.stringify(change._id));
});
}
startOrderListener().catch(console.error);
// 2. Manual Testing in mongosh
// Trigger Event:Insert a New Order
db.orders.insertOne({
userId: 'user_001',
items: [{ sku: 'PHONE-001', qty: 1, price: 599 }],
total: 599,
status: 'paid',
createdAt: new Date()
});
// The app receives a notification immediately:Send an Email + WebSocket Push + ES Index
الناتج: يتم تشغيله في غضون أقل من 100 مللي ثانية بعد إدراج الأمر، ويقوم تلقائيًا بإتمام جميع الآثار الجانبية — مثل رسائل البريد الإلكتروني والإشعارات الفورية ومزامنة البحث — دون الحاجة إلى إجراء استقصاء.
▶ مثال 3: مراقبة تغييرات المخزون في الوقت الفعلي مع التنبيهات(الصعوبة ⭐⭐)
// نظام TechCorp لمراقبة المخزون: تنبيه فوري عند نفاد المخزون
const Product = mongoose.model('Product', new mongoose.Schema({
sku: String,
title: String,
stock: Number,
lowStockThreshold: { type: Number, default: 10 }
}));
// مراقبة تحديثات المخزون
const stockStream = Product.watch([
{
$match: {
operationType: 'update',
'updateDescription.updatedFields.stock': { $exists: true }
}
}
], { fullDocument: 'updateLookup' });
stockStream.on('change', async (change) => {
const product = change.fullDocument;
if (product.stock <= product.lowStockThreshold) {
// إرسال تنبيه للمخزون المنخفض
console.log(`⚠️ تحذير: المنتج ${product.sku} - ${product.title}`);
console.log(` المخزون الحالي: ${product.stock}`);
console.log(` الحد الأدنى: ${product.lowStockThreshold}`);
// إرسال بريد إلكتروني للمشرف
await sendEmail('inventory@techcorp.com', {
subject: `تحذير: مخزون منخفض - ${product.sku}`,
body: `المنتج "${product.title}" وصل إلى ${product.stock} وحدة فقط`
});
// إرسال إشعار فوري للوحة التحكم
io.emit('low_stock_alert', {
sku: product.sku,
title: product.title,
currentStock: product.stock,
threshold: product.lowStockThreshold
});
}
});
// محاكاة تحديث المخزون
await Product.updateOne(
{ sku: 'PHONE-001' },
{ $inc: { stock: -50 } }
);
الإخراج:
TEXT 📖 للعرض فقط⚠️ تحذير: المنتج PHONE-001 - هاتف ذكي X المخزون الحالي: 5 الحد الأدنى: 10 [تم إرسال بريد إلكتروني إلى inventory@techcorp.com]
❓ أسئلة شائعة
resumeAfter لاستئناف العملية من حيث توقفت.resumeToken لتخزينه خارجيًا (على سبيل المثال، في Redis) حتى يمكن استعادته عند إعادة التشغيل.📖 ملخص
- تيارات التغييرات: إخطارات بالتغييرات في الوقت الفعلي استنادًا إلى سجل العمليات (oplog)
- طريقة watch() للاستماع + تصفية $match
- نوع العملية: إدراج / تحديث / حذف / استبدال / حذف الكائن
- fullDocument: 'updateLookup' الحصول على المستند الكامل
- السيناريوهات النموذجية: الإشعارات في الوقت الفعلي، مركز السيطرة على الأمراض (CDC)، إبطال صلاحية ذاكرة التخزين المؤقت
- مجموعات النسخ المتماثلة المطلوبة، الحد الأقصى لحجم سجل العمليات (oplog)
📝 تمارين
- المشكلة الأساسية (⭐): مراقبة جميع التغييرات التي تطرأ على المجموعة
productsوعرضها على وحدة التحكم. - السؤال الأساسي (⭐): استخدم $match لتصفية عمليات الإدراج.
- تمرين متقدم (⭐⭐): استمع إلى الطلبات الجديدة وأرسل إشعارات عبر البريد الإلكتروني (باستخدام وظيفة بريد إلكتروني وهمية).
- تمرين متقدم (⭐⭐): تنفيذ انتهاء صلاحية ذاكرة التخزين المؤقت للمنتجات (مزامنة Redis).
- التحدي (⭐⭐⭐): إكمال نظام CDC (المزامنة في الوقت الفعلي بين MongoDB وElasticsearch).