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

使用Axios批量异步调用持续等待响应问题排查求助

问题排查与修复:Node.js批量API调用无响应问题

问题背景

运行Node.js代码通过Axios批量异步调用外部API时,请求已送达服务器并得到响应,但代码始终无法接收响应,导致批量处理停滞。Postman使用相同参数测试正常,推测问题与Async/Await及流处理逻辑相关,需求是处理大型CSV文件,分批并发发送请求。

核心问题分析

  1. 流事件的异步处理未控制背压
    data事件的异步回调不会被stream等待,CSV流会持续快速推送数据,导致批量逻辑混乱,请求上下文丢失,后续响应无法被正确捕获。

  2. Axios请求配置错误
    手动使用JSON.stringify(data)作为请求体,但Axios会自动序列化JSON,这会导致请求体格式异常,服务器响应后客户端解析失败。

  3. 工作池与队列逻辑缺陷

    • 任务入队后没有主动触发队列消费,仅在end事件中调用processQueue,可能导致队列任务堆积。
    • processQueue采用递归调用,大量任务时可能导致栈溢出,且未正确等待异步任务完成后再处理下一个队列任务。
  4. 全局变量不安全
    configList、configProcessed等全局变量在多请求场景下会被污染,且初始化逻辑存在风险。

  5. Promise错误处理不当
    Promise.all会因单个请求失败而终止整个批量,且callAPI中返回的Promise.reject未被正确捕获处理。

修复方案

1. 修复流处理的异步控制

使用through2的异步回调支持,在转换流中控制数据处理节奏,避免背压;改用局部变量存储批量数据,替换全局变量。

2. 修正Axios请求配置

移除JSON.stringify(data),直接传递data对象,让Axios自动处理JSON序列化。

3. 完善工作池与队列逻辑

  • 在工作线程释放后触发队列消费,确保入队任务能被及时处理。
  • 将processQueue改为循环异步处理,避免递归栈溢出。

4. 替换全局变量为局部变量

在bootstrapConfig内部定义批量数据和计数变量,避免多请求冲突。

5. 改进错误处理

使用Promise.allSettled替代Promise.all,确保批量中部分请求失败不影响整体流程,同时记录每个请求的结果。

完整修复代码

var express = require('express');
var router = express.Router();

const fs = require('fs');
const csv = require('csv-stream');
const through2 = require('through2');
const axios = require('axios');

// 配置参数(建议从环境变量或配置文件读取)
const MAX_CONCURRENT_BATCHES = 5;
const BATCH_SIZE = 10; // 根据实际需求调整
const API_URL = '你的API地址';
const BASIC_AUTH = '你的Base64认证字符串';

// 工作池:跟踪当前运行的批量任务数
let activeBatches = 0;
const taskQueue = [];

// 处理队列任务的异步循环
async function processQueue() {
    while (taskQueue.length > 0 && activeBatches < MAX_CONCURRENT_BATCHES) {
        activeBatches++;
        const taskData = taskQueue.shift();
        try {
            await processConfigurations(taskData);
        } catch (error) {
            console.error('批量处理失败:', error);
        } finally {
            activeBatches--;
        }
    }
}

// 调用单个API
async function callAPI(data) {
    try {
        const response = await axios.request({
            method: 'PUT',
            url: API_URL,
            data: data, // 直接传对象,Axios自动序列化
            headers: {
                'Content-Type': 'application/json',
                'Authorization': 'Basic ' + BASIC_AUTH
            }
        });
        return { success: true, data: response.data };
    } catch (error) {
        let message = '未知错误: ' + error.message;
        if (error.response) {
            message = `HTTP ${error.response.status}: ${JSON.stringify(error.response.data)}`;
        }
        console.error('API调用失败:', message);
        return { success: false, error: message };
    }
}

// 处理单批请求
async function processConfigurations(taskData) {
    const results = await Promise.allSettled(
        taskData.configList.map(config => callAPI(config))
    );
    // 统计处理结果
    const processed = results.length;
    const successCount = results.filter(r => r.status === 'fulfilled' && r.value.success).length;
    console.log(`批量处理完成: 总请求${processed},成功${successCount}`);
    return results;
}

// 提交任务到队列或直接执行
async function submitTask(taskData) {
    if (activeBatches < MAX_CONCURRENT_BATCHES) {
        activeBatches++;
        try {
            await processConfigurations(taskData);
        } catch (error) {
            console.error('任务执行失败:', error);
        } finally {
            activeBatches--;
            processQueue(); // 任务完成后触发队列处理
        }
    } else {
        taskQueue.push(taskData);
    }
}

// 处理CSV文件
async function bootstrapConfig(fileName) {
    let configList = [];
    let totalProcessed = 0;

    return new Promise((resolve, reject) => {
        const stream = fs.createReadStream(`./${fileName}.csv`)
            .pipe(csv.createStream({
                endLine: '\n',
                columns: ['COL1', 'COL2', 'COL3', 'COL4'],
                escapeChar: '"',
                enclosedChar: '"'
            }))
            .pipe(through2({ objectMode: true }, async (row, enc, cb) => {
                // 构造请求数据(根据实际需求修改)
                const config = {
                    col1: row.COL1,
                    col2: row.COL2,
                    col3: row.COL3,
                    col4: row.COL4
                };
                configList.push(config);
                totalProcessed++;

                // 达到批量大小则提交任务
                if (configList.length === BATCH_SIZE) {
                    const batch = [...configList];
                    configList = [];
                    await submitTask({ configList: batch });
                }
                cb();
            }))
            .on('end', async () => {
                // 处理剩余的不足批量的数据
                if (configList.length > 0) {
                    await submitTask({ configList: configList });
                }
                // 等待队列中所有任务完成
                while (activeBatches > 0 || taskQueue.length > 0) {
                    await new Promise(resolve => setTimeout(resolve, 100));
                }
                console.log(`所有任务处理完成,总记录数: ${totalProcessed}`);
                resolve(totalProcessed);
            })
            .on('error', err => {
                console.error('CSV读取错误:', err);
                reject(err);
            });
    });
}

// API接口
router.get('/', async function(req, res, next) {
    try {
        // 假设fileName从请求参数获取,或固定值
        const fileName = req.query.fileName || 'data';
        await bootstrapConfig(fileName);
        res.json({ status: "Success", message: "所有任务已启动处理,可查看日志跟踪进度" });
    } catch (error) {
        res.status(500).json({ status: "Error", message: error.message });
    }
});

module.exports = router;

关键修复说明

  • 流处理控制:使用through2的异步回调,确保每一行处理完成后再继续,避免数据堆积。
  • Axios配置修正:移除手动JSON序列化,让Axios自动处理请求体格式。
  • 队列与工作池优化:用异步循环处理队列,避免递归栈溢出,任务完成后自动触发下一个任务。
  • 局部变量替代全局变量:避免多请求场景下的变量污染。
  • 错误处理增强:Promise.allSettled确保批量中部分失败不中断整体流程,同时记录每个请求的结果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 10:24:53