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
相关产品推荐
相关产品推荐

