Skip to content

Commit f45a1e0

Browse files
committed
feat: add watch() method with complete documentation and examples (v1.1.0)
- Add watch() method for MongoDB Change Streams - Auto-reconnect with exponential backoff - Smart cache invalidation - Cross-instance synchronization - Complete documentation and README examples - 24 tests (17 unit + 7 integration) - Replica set support
1 parent a4fe879 commit f45a1e0

13 files changed

Lines changed: 2383 additions & 24 deletions

File tree

CHANGELOG.md

Lines changed: 41 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -2,18 +2,47 @@
22

33
所有显著变更将记录在此文件,遵循 Keep a Changelog 与语义化版本(SemVer)。
44

5-
## [Unreleased] - v1.1.0 计划
6-
7-
### 🗺️ 计划功能
8-
9-
**Change Streams(实时监听)**
10-
- watch API - 监听集合/数据库变更
11-
- 智能缓存失效 - watch 事件自动失效相关缓存
12-
- 跨实例缓存同步 - 分布式环境下的缓存一致性
13-
- 自动重连机制 - 网络中断后自动恢复
14-
- 断点续传 - resumeToken 自动保存和恢复
15-
16-
**预计发布**: 2025-12 下旬
5+
## [1.1.0] - 2025-12-03
6+
7+
### 🎊 v1.1.0 发布 - Change Streams 支持
8+
9+
**新功能**:
10+
-**watch() 方法** - MongoDB Change Streams 实时监听
11+
- 监听集合/数据库的数据变更(insert/update/delete/replace)
12+
- 支持聚合管道过滤事件
13+
- 完整的事件系统(change, error, reconnect, resume, close, fatal)
14+
15+
- 🔄 **自动重连机制**
16+
- 网络中断后自动重连(指数退避算法:1s → 2s → 4s → 8s → ... → 60s)
17+
- resumeToken 自动管理和断点续传
18+
- 智能错误分类(瞬态/持久性/致命)
19+
20+
- 🗑️ **智能缓存失效**
21+
- 监听到数据变更时自动失效相关缓存
22+
- 支持精准失效(根据 operationType 和 documentKey)
23+
- 自动触发跨实例缓存同步(复用 DistributedCacheInvalidator)
24+
25+
- 📊 **统计监控**
26+
- getStats() 方法获取运行统计
27+
- 监控总变更数、重连次数、缓存失效次数等
28+
29+
- 🧪 **副本集支持**
30+
- mongodb-memory-server 副本集模式配置
31+
- 测试环境完整支持 Change Streams
32+
33+
**文档**:
34+
- 📄 新增 `docs/watch.md` - 完整 API 文档(400 行)
35+
- 📝 新增 `examples/watch.examples.js` - 6 个使用示例
36+
- 🔗 更新 `docs/events.md``docs/INDEX.md` - 交叉引用
37+
38+
**测试**:
39+
- ✅ 新增 17 个单元测试(100% 通过)
40+
- ✅ 新增 7 个集成测试(100% 通过,副本集环境)
41+
- ✅ 测试覆盖率约 85%
42+
43+
**兼容性**:
44+
- ✅ 完全向后兼容
45+
- ✅ 纯新增功能,无破坏性变更
1746

1847
---
1948

README.md

Lines changed: 136 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -270,6 +270,7 @@ console.log(`总计: ${result.totals?.total}, 共 ${result.totals?.totalPages}
270270
-**Read**: find, findOne, findPage(游标分页), aggregate, count, distinct
271271
-**Update**: updateOne, updateMany, replaceOne, findOneAndUpdate, findOneAndReplace
272272
-**Delete**: deleteOne, deleteMany, findOneAndDelete
273+
-**Watch**: watch(Change Streams 实时监听)**⭐ v1.1.0**
273274
274275
#### **索引管理(100% 完成)**
275276
- ✅ createIndex, createIndexes, listIndexes, dropIndex, dropIndexes
@@ -980,6 +981,141 @@ const salesReport = await collection.aggregate([
980981
981982
---
982983
984+
### 实时监听(watch)⭐ v1.1.0
985+
986+
**监听 MongoDB 数据变更,支持自动缓存失效**
987+
988+
#### 1. 基础监听
989+
990+
```javascript
991+
// 监听集合的所有数据变更
992+
const watcher = collection.watch();
993+
994+
watcher.on('change', (change) => {
995+
console.log('数据变更:', change.operationType); // insert/update/delete/replace
996+
console.log('文档ID:', change.documentKey._id);
997+
console.log('完整文档:', change.fullDocument);
998+
});
999+
1000+
// 插入数据(会触发 change 事件)
1001+
await collection.insertOne({ name: 'Alice', age: 25 });
1002+
```
1003+
1004+
#### 2. 过滤事件
1005+
1006+
```javascript
1007+
// 只监听 insert 和 update 操作
1008+
const watcher = collection.watch([
1009+
{ $match: { operationType: { $in: ['insert', 'update'] } } }
1010+
]);
1011+
1012+
watcher.on('change', (change) => {
1013+
console.log('新增或修改:', change.operationType);
1014+
});
1015+
```
1016+
1017+
#### 3. 自动缓存失效 ⭐
1018+
1019+
```javascript
1020+
// 启用自动缓存失效(默认开启)
1021+
const watcher = collection.watch([], {
1022+
autoInvalidateCache: true // 数据变更时自动失效相关缓存
1023+
});
1024+
1025+
// 1. 查询并缓存数据
1026+
const users = await collection.find({ status: 'active' }, { cache: 60000 });
1027+
1028+
// 2. 更新数据(触发 watch)
1029+
await collection.updateOne({ _id: userId }, { $set: { status: 'inactive' } });
1030+
1031+
// 3. ✅ watch 自动失效相关缓存
1032+
// 4. 下次查询自动从数据库读取最新数据
1033+
```
1034+
1035+
#### 4. 错误处理和重连
1036+
1037+
```javascript
1038+
const watcher = collection.watch();
1039+
1040+
// 监听错误(自动重试瞬态错误)
1041+
watcher.on('error', (error) => {
1042+
console.warn('持久性错误:', error.message);
1043+
});
1044+
1045+
// 监听重连
1046+
watcher.on('reconnect', (info) => {
1047+
console.log(`第 ${info.attempt} 次重连,延迟 ${info.delay}ms`);
1048+
});
1049+
1050+
// 监听恢复
1051+
watcher.on('resume', () => {
1052+
console.log('✅ 已恢复监听(断点续传)');
1053+
});
1054+
1055+
// 监听致命错误
1056+
watcher.on('fatal', (error) => {
1057+
console.error('💥 致命错误(无法恢复):', error);
1058+
// 通知运维
1059+
});
1060+
```
1061+
1062+
#### 5. 统计监控
1063+
1064+
```javascript
1065+
const watcher = collection.watch();
1066+
1067+
// 获取运行统计
1068+
const stats = watcher.getStats();
1069+
console.log('总变更数:', stats.totalChanges);
1070+
console.log('重连次数:', stats.reconnectAttempts);
1071+
console.log('运行时长:', stats.uptime, 'ms');
1072+
console.log('缓存失效次数:', stats.cacheInvalidations);
1073+
console.log('活跃状态:', stats.isActive);
1074+
```
1075+
1076+
#### 6. 优雅关闭
1077+
1078+
```javascript
1079+
// 应用退出时关闭 watcher
1080+
process.on('SIGTERM', async () => {
1081+
await watcher.close();
1082+
await db.close();
1083+
process.exit(0);
1084+
});
1085+
```
1086+
1087+
**核心特性**
1088+
-**自动重连**:网络中断后自动恢复(指数退避:1s2s4s...60s
1089+
-**断点续传**:resumeToken 自动管理,不丢失任何变更
1090+
-**智能缓存失效**:数据变更时自动失效相关缓存
1091+
-**跨实例同步**:分布式环境自动广播缓存失效
1092+
-**完整事件系统**:change, error, reconnect, resume, close, fatal
1093+
-**统计监控**:完整的运行统计和健康检查
1094+
1095+
**注意事项**
1096+
- ⚠️ **需要副本集**:Change Streams 需要 MongoDB 4.0+ 副本集或分片集群
1097+
- ⚠️ **测试环境**:可使用 mongodb-memory-server 副本集模式
1098+
1099+
**测试环境配置**
1100+
```javascript
1101+
const db = new MonSQLize({
1102+
type: 'mongodb',
1103+
databaseName: 'mydb',
1104+
config: {
1105+
useMemoryServer: true,
1106+
memoryServerOptions: {
1107+
instance: {
1108+
replSet: 'rs0' // 启用副本集(支持 Change Streams)
1109+
}
1110+
}
1111+
}
1112+
});
1113+
```
1114+
1115+
📖 详细文档:[watch 方法完整指南](./docs/watch.md)
1116+
1117+
---
1118+
9831119
## 📚 完整文档
9841120
9851121
### 核心文档

STATUS.md

Lines changed: 11 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -54,17 +54,18 @@
5454
| **MongoDB 事务** | 8 | 0 | 0 | 0 | 8 |
5555
| **分布式支持** | 3 | 0 | 0 | 0 | 3 |
5656
| **MongoDB Admin/Management** | 18 | 0 | 0 | 0 | 18 |
57-
| **MongoDB Change Streams** | 0 | 1 | 0 | 0 | 1 |
57+
| **MongoDB Change Streams** | 1 | 0 | 0 | 0 | 1 |
5858
| **MongoDB 其他** | 0 | 0 | 0 | 1 | 1 |
59-
| **总计** | **89** | **1** | **2** | **2** | **94** |
59+
| **总计** | **90** | **0** | **2** | **2** | **94** |
6060

61-
**完成度**: **100%** (89/89,不含计划中、不推荐和其他功能) 🎉
61+
**完成度**: **100%** (90/90,不含计划中、不推荐和其他功能) 🎉
6262
**核心功能完成度**: **100%** (30/30) 🎊
6363
**CRUD + 索引 + 事务 + 便利性方法完成度**: **100%**
6464
**Admin/Management 功能完成度**: **100%** (18/18) ✅
65+
**Change Streams 功能完成度**: **100%** (1/1) ✅
6566
**文档完成度**: **95%+** (findAndCount 文档待补充) ✅
6667

67-
**v1.1.0 计划**: watch(Change Streams)实时监听功能
68+
**v1.1.0 已发布**: watch(Change Streams)实时监听功能
6869

6970
### 📈 进度对比
7071

@@ -80,12 +81,13 @@
8081
| 2025-12-02 (上午) | 93.3% | **✨ 便利方法完成** (4/5) + 文档改进 + STATUS 优化 🎉 |
8182
| 2025-12-02 (下午) | **98.9%** | **🛠️ Admin/Management 完成** (18个方法,102个测试100%通过) 🎉 |
8283
| 2025-12-03 | **100%** | **🎉 v1.0.0 正式版发布** - 已成功发布到 npm! |
83-
| 增长 | **+1.1%** | v1.0.0 稳定发布,企业级质量标准达成 |
84+
| 2025-12-03 (下午) | **100%** | **🎊 v1.1.0 watch 功能完成** - Change Streams 实时监听 + 自动缓存失效 🎉 |
85+
| 增长 | **稳定** | v1.1.0 watch 功能完整实现,企业级质量标准达成 |
8486

85-
**v1.1.0 路线图**:
86-
| 计划日期 | 目标功能 | 说明 |
87-
|---------|---------|------|
88-
| 2025-12 下旬 | **watch (Change Streams)** | 实时监听功能,支持智能缓存失效、跨实例同步 |
87+
**v1.1.0 已发布**:
88+
| 功能 | 状态 | 说明 |
89+
|------|------|------|
90+
| **watch (Change Streams)** | ✅ 已完成 | 实时监听功能,支持自动重连、断点续传、智能缓存失效、跨实例同步 |
8991

9092
### 📚 2025-11-18 文档补全成果 (阶段1)
9193

docs/INDEX.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
| [findPage.md](findPage.md) | `findPage()` | 游标分页查询 |
3030
| [count.md](count.md) | `count()` | 统计文档数量 |
3131
| [distinct.md](distinct.md) | `distinct()` | 去重查询 |
32+
| [watch.md](watch.md) | `watch()` | **实时监听数据变更(Change Streams)⭐** |
3233

3334
---
3435

docs/events.md

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -723,6 +723,34 @@ msq2.on('slow-query', () => console.log('msq2 慢查询'));
723723

724724
---
725725

726+
## 相关文档
727+
728+
### watch 事件
729+
730+
如果你需要监听 MongoDB 的数据变更(而非应用的查询操作),请参考:
731+
732+
- [watch 方法文档](./watch.md) - MongoDB Change Streams
733+
734+
**区别**:
735+
- 全局事件(本文档):监听应用的查询操作
736+
- watch 事件:监听 MongoDB 的数据变更
737+
738+
**示例**:
739+
```javascript
740+
// 全局事件:监听应用的慢查询
741+
msq.on('slow-query', (meta) => {
742+
console.warn('应用执行了慢查询');
743+
});
744+
745+
// watch 事件:监听 MongoDB 的数据变更
746+
const watcher = collection.watch();
747+
watcher.on('change', (change) => {
748+
console.log('MongoDB 数据变更');
749+
});
750+
```
751+
752+
---
753+
726754
## 参考资料
727755

728756
- [Node.js EventEmitter 文档](https://nodejs.org/api/events.html)

0 commit comments

Comments
 (0)