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

Node.js中Express流式x-ndjson测试阻塞与Promise未解析问题排查

问题描述

我有一个使用Node v19.1.0的TypeScript库,该库包含一个监听流式服务器事件的函数。服务器提供/events路由,流式传输'application/x-ndjson'内容,类型可能为event/ping等(需定时发送ping维持连接)。我的observe函数解析流式数据,将有效事件传递给回调,调用者可通过中止函数按需终止流。

本地或CI运行测试时,会出现以下错误:

Warning: Test "observes events." generated asynchronous activity after the test ended. This activity created the error "AbortError: The operation was aborted." and would have caused the test to fail, but instead triggered an unhandledRejection event.

我用纯JavaScript简化了示例代码:

const assert = require('assert/strict');
const express = require('express');
const { it } = require('node:test');

it('observes events.', async () => {
    const expectedEvent = { type: 'event', payload: { metadata: { type: 'entity-created', commandId: 'commandId' } } };

    const api = express();
    const server = api
        .use(express.json())
        .post('/events', (request, response) => {
            response.writeHead(200, {
                'content-type': 'application/x-ndjson',
            });

            const line = JSON.stringify(expectedEvent) + '\n';
            response.write(line);
        })
        .listen(3000);

    let stopObserving = () => {
        throw new Error('should never happen');
    };

    const actualEventPayload = await new Promise(async resolve => {
        stopObserving = await observeEvents(async newEvent => {
            resolve(newEvent);
        });
    });

    stopObserving();
    server.closeAllConnections();
    server.close();

    assert.deepEqual(actualEventPayload, expectedEvent.payload);
});

const observeEvents = async function (onReceivedFn) {
    const abortController = new AbortController();

    const response = await fetch('http://localhost:3000/events', {
        method: 'POST',
        headers: { 'content-type': 'application/json' },
        signal: abortController.signal,
    });

    if (!response.ok) {
        throw new Error('error handling goes here - request failed');
    }

    Promise.resolve().then(async () => {
        if (!response.body) {
            throw new Error('error handling goes here - missing response body');
        }

        for await (const item of parseStream(response.body, abortController)) {
            switch (item.type) {
                case 'event': {
                    await onReceivedFn(item.payload);

                    break;
                }
                case 'ping':
                    // Intentionally left blank
                    break;
                case 'error':
                    throw new Error('error handling goes here - stream failed');
                default:
                    throw new Error('error handling goes here - should never happen');
            }
        }
    });

    return () => { abortController.abort(); };
};

const parseLine = function () {
    return new TransformStream({
        transform(chunk, controller) {
            try {
                const data = JSON.parse(chunk);

                // ... check if this is a valid line...

                controller.enqueue(data);
            } catch (error) {
                controller.error(error);
            }
        },
    });
};

const splitLines = function () {
    let buffer = '';

    return new TransformStream({
        transform(chunk, controller) {
            buffer += chunk;

            const lines = buffer.split('\n');

            for (let i = 0; i < lines.length - 1; i++) {
                controller.enqueue(lines[i]);
            }

            buffer = lines.at(-1) ?? '';
        },
        flush(controller) {
            if (buffer.length > 0) {
                controller.enqueue(buffer);
            }
        },
    });
};

const parseStream = async function* (stream, abortController) {
    let streamReader;

    try {
        const pipedStream = stream
            .pipeThrough(new TextDecoderStream())
            .pipeThrough(splitLines())
            .pipeThrough(parseLine());

        streamReader = pipedStream.getReader();

        while (true) {
            const item = await streamReader.read();

            if (item.done) {
                break;
            }

            yield item.value;
        }
    } finally {
        await streamReader?.cancel();
        abortController.abort();
    }
};

运行node --test时测试无法结束,必须手动取消,问题卡在以下代码段:

const actualEventPayload = await new Promise(async resolve => {
    stopObserving = await observeEvents(async newEvent => {
        resolve(newEvent);
    });
});

我认为是Promise从未解析导致的,即使移除所有流解析代码,替换相关逻辑后问题仍存在,请问问题出在哪里或缺少什么?


问题分析与解决

核心问题1:服务器响应未结束,流长期挂起

测试服务器的/events路由仅调用response.write(line)发送数据,但未调用response.end()结束响应。这会导致客户端的fetch请求一直处于等待状态,流无法触发done信号,parseStream中的循环会无限卡住。

核心问题2:异步任务未处理错误,触发未捕获rejection

observeEvents中用Promise.resolve().then(async () => { ... })启动了独立异步任务,但未处理该任务的错误。当调用abortController.abort()时,流的read()方法会抛出AbortError,这个错误未被捕获,就会触发测试结束后的异步活动警告。

核心问题3:服务器关闭逻辑未等待完成

直接调用server.closeAllConnections()和server.close()但未等待关闭完成的回调,测试可能提前结束,但服务器连接或进程仍在运行,残留异步活动。

修复步骤

  1. 修改服务器路由,结束响应
    在response.write(line)后添加response.end(),告知客户端流已结束:

    .post('/events', (request, response) => {
        response.writeHead(200, {
            'content-type': 'application/x-ndjson',
        });
    
        const line = JSON.stringify(expectedEvent) + '\n';
        response.write(line);
        response.end(); // 新增:结束响应
    })
    
  2. 处理observeEvents中的异步任务错误
    移除独立的Promise.resolve().then()链,直接在函数内部处理流逻辑,并捕获错误:

    const observeEvents = async function (onReceivedFn) {
        const abortController = new AbortController();
    
        const response = await fetch('http://localhost:3000/events', {
            method: 'POST',
            headers: { 'content-type': 'application/json' },
            signal: abortController.signal,
        });
    
        if (!response.ok) {
            throw new Error('error handling goes here - request failed');
        }
    
        if (!response.body) {
            throw new Error('error handling goes here - missing response body');
        }
    
        try {
            for await (const item of parseStream(response.body, abortController)) {
                switch (item.type) {
                    case 'event': {
                        await onReceivedFn(item.payload);
                        break;
                    }
                    case 'ping':
                        break;
                    case 'error':
                        throw new Error('error handling goes here - stream failed');
                    default:
                        throw new Error('error handling goes here - should never happen');
                }
            }
        } catch (err) {
            // 忽略AbortError,其他错误重新抛出
            if (err.name !== 'AbortError') {
                throw err;
            }
        }
    
        return () => { abortController.abort(); };
    };
    
  3. 优化测试中的服务器关闭逻辑
    改为异步等待服务器关闭完成,避免残留异步活动:

    stopObserving();
    await new Promise((resolve) => {
        server.closeAllConnections();
        server.close(resolve); // 等待服务器关闭完成
    });
    
  4. 简化Promise初始化逻辑
    不需要在new Promise中使用async resolve,直接调整为:

    let stopObserving;
    const actualEventPayload = await new Promise(resolve => {
        observeEvents(async newEvent => {
            resolve(newEvent);
        }).then(fn => stopObserving = fn);
    });
    

修复后流程说明

  • 服务器发送事件后立即结束响应,客户端流处理完数据后正常退出循环
  • observeEvents中的异步错误被捕获,不会触发未处理的rejection
  • 测试等待服务器完全关闭后再结束,彻底消除残留异步活动

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 02:15:33