使用Observable.create实现defer式Observable的Node.js执行异常排查
让我一步步帮你拆解问题,找出你代码里的坑,以及为什么会出现这些奇怪的现象:
首先,你犯了一个典型的异步写法错误
看你代码里的这部分:
fetch('https://jsonplaceholder.typicode.com/todos/1').then(console.log('fetch done'))
这里的console.log('fetch done')是立即执行的,根本不是等到请求完成后才打印!因为then()方法需要接收一个函数作为参数,而你直接把console.log('fetch done')的执行结果(也就是undefined)传给了then,相当于没设置任何回调。所以不管请求成功还是失败,这行代码一执行就会打印fetch done,这和你预期的“请求完成后打印”完全不符。
正确的写法应该是把console.log包装成一个函数:
fetch('...').then(() => console.log('fetch done'))
其次,你没搞懂Observable.create和defer的核心差异
你想模仿defer的行为——每次订阅时才触发新的异步操作,但你的写法没抓住重点:
第一个版本(无setTimeout)的问题
当你调用obs.subscribe()时,Observable.create的回调才会执行,此时会立即发起fetch请求,然后把Promise转换成Observable并绑定到observer。但为什么程序会无限等待?
- 你只给
subscribe传了next回调,没处理complete和error事件。在Node.js中,如果RxJS的Observable没有触发complete或error,进程可能会因为残留的异步订阅而无法退出。 - 更关键的是,你用
from(fetch(...)).subscribe(observer)的方式绑定订阅,但没返回清理函数。当外部订阅取消时,内部的订阅不会自动清理,可能导致进程挂起。
第二个版本(用setTimeout包裹)的问题
你把内部Observable的创建放到了setTimeout(() => ..., 0)里,这导致请求逻辑被推迟到下一个宏任务队列。为什么有时成功有时失败?
- 竞态条件:如果fetch请求的Promise回调和其他宏任务的执行顺序不确定,可能会出现内部Observable还没发出值,外部订阅就已经处于未就绪状态;
- 未处理的错误:如果请求失败(比如网络波动),你没处理
error回调,RxJS会抛出未捕获的异常,导致进程异常挂起,这时候你只会看到提前打印的fetch done,看不到OK。
正确的写法(模仿defer的效果)
如果你非要用Observable.create实现defer的行为,需要做到这几点:
- 每次订阅时才发起请求(这一点你第一个版本其实做到了,但被then的错误写法掩盖了)
- 返回清理函数,确保外部订阅取消时,内部订阅也能被清理
- 正确处理Promise的回调和错误
代码如下:
const obs = Observable.create(function(observer) { // 每次订阅才发起请求 const fetchPromise = fetch('https://jsonplaceholder.typicode.com/todos/1') .then(resp => { console.log('fetch done'); return resp; // 把响应传给Observable }) .catch(err => { observer.error(err); // 把错误传给订阅者 throw err; }); // 绑定内部Observable到外部observer const subscription = from(fetchPromise).subscribe(observer); // 返回清理函数:外部订阅取消时,取消内部订阅 return () => { subscription.unsubscribe(); // 如果需要取消请求,可以用AbortController // controller.abort(); }; }); // 5秒后订阅,同时处理next/error/complete setTimeout(() => { obs.subscribe( (resp) => console.log(resp.statusText), (err) => console.error('请求出错:', err), () => console.log('请求完成') // 触发complete,确保进程能正常退出 ); }, 5000);
当然,如果你只是想实现defer的效果,直接用RxJS的defer操作符会更简单,不用自己造轮子:
const obs = defer(() => fetch('https://jsonplaceholder.typicode.com/todos/1') .then(resp => { console.log('fetch done'); return resp; }) ); setTimeout(() => { obs.subscribe( (resp) => console.log(resp.statusText), (err) => console.error('请求出错:', err), () => console.log('请求完成') ); }, 5000);
总结你踩的坑
- 异步回调写法错误:直接执行
console.log传给then,导致打印时机完全不符合预期。 - 订阅生命周期处理不当:用
Observable.create时没返回清理函数,可能导致订阅残留,进程无法退出。 - 缺少错误和complete处理:订阅时只关注
next事件,忽略了错误和完成事件,导致进程异常挂起。
内容的提问来源于stack exchange,提问作者Marinos An
相关产品推荐
相关产品推荐

