You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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的行为,需要做到这几点:

  1. 每次订阅时才发起请求(这一点你第一个版本其实做到了,但被then的错误写法掩盖了)
  2. 返回清理函数,确保外部订阅取消时,内部订阅也能被清理
  3. 正确处理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);

总结你踩的坑

  1. 异步回调写法错误:直接执行console.log传给then,导致打印时机完全不符合预期。
  2. 订阅生命周期处理不当:用Observable.create时没返回清理函数,可能导致订阅残留,进程无法退出。
  3. 缺少错误和complete处理:订阅时只关注next事件,忽略了错误和完成事件,导致进程异常挂起。

内容的提问来源于stack exchange,提问作者Marinos An

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.06 17:59:10