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

RxJS实现TCP服务器时Observable自递归问题求助

问题原因与修复方案

核心错误1:错误的发送方式触发无限循环

你在输出订阅的next回调里用了socket.emit("data", response),这是触发socket自身的data事件——相当于把要发给客户端的消息,又当成客户端发来的输入被input$捕获,接着再次处理、发送,形成无限循环,这就是为什么会出现数千次重复ACK的原因。

正确的TCP客户端数据发送方式是使用socket.write(response),它只会把数据写入socket发送给客户端,不会触发本地的data事件。

核心错误2:重复订阅input$导致多次监听

你的handler里先通过await rx.firstValueFrom(input$)订阅了一次input$,之后返回的Observable又被订阅了一次。而你创建的input$每次订阅都会给socket添加新的data事件监听器,这意味着客户端的同一条消息会被多个监听器捕获,进一步放大了循环次数。

修复后的代码

调整发送方式,同时将登录逻辑改为流式处理,避免重复订阅:

import net from "net";
import * as rx from "rxjs";
import { first, map, switchMap, tap } from "rxjs/operators";

type MaybePromise<T> = T | PromiseLike<T>;
type ConnectionHandler = (address: net.AddressInfo, input$: rx.Observable<Buffer>) => MaybePromise<rx.Observable<Buffer>>;

const createServer = (port: number, handler: ConnectionHandler, opts?: net.ServerOpts) => {
  const server = net.createServer(opts);

  new rx.Observable<net.Socket>(subscriber => {
    server
      .on("connection", socket => subscriber.next(socket))
      .on("error", err => subscriber.error(err));
  }).subscribe(async socket =>
    (await handler(
      socket.address() as net.SocketAddress,
      new rx.Observable<Buffer>(subscriber => {
        const onData = (request: Buffer) => subscriber.next(request);
        const onError = (err: Error) => subscriber.error(err);
        const onClose = () => subscriber.complete();
        
        socket
          .on("end", () => socket.destroySoon())
          .on("error", onError)
          .on("data", onData)
          .on("close", onClose);
        
        // 订阅取消时移除事件监听,避免内存泄漏
        return () => {
          socket.off("error", onError);
          socket.off("data", onData);
          socket.off("close", onClose);
        };
      }),
    )).subscribe({
      next: response => {
        console.log(response.toString());
        // 替换为正确的客户端发送方式
        socket.write(response);
      },
      error: () => socket.destroySoon(),
    })
  );

  server.listen(port);
};

createServer(3000, (address, input$) => {
  console.log(`New client: ${address.address}`);

  return input$.pipe(
    // 取第一条消息做登录验证
    first(),
    tap(login => {
      if (login.toString() !== "login\n") {
        throw new Error("login failure");
      }
    }),
    // 验证通过后切换到后续消息流,跳过已处理的登录消息
    switchMap(() => input$.pipe(
      map(data => Buffer.from(`ACK: ${data.toString()}`, "utf-8"))
    ))
  );
});

额外优化说明

  • 给input$的Observable添加了清理逻辑,在订阅取消时移除socket的事件监听器,避免内存泄漏。
  • 登录逻辑改为流式操作,整个过程只订阅一次input$,避免重复监听socket事件。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 15:32:41