Node.js应用运行2小时后出现EADDRNOTAVAIL错误求助
问题描述
我开发了一个从Kafka消费数据并发送至名为DB的容器的Node.js应用。程序启动后可正常运行约2小时,随后消费者端的Axios捕获到如下错误:
Api error :connect:EADDRNOTAVAL:3002(该端口为DB容器的端口)
补充说明:报错时DB端无任何日志或错误输出,主机环境为Linux CentOS 8。
消费者端代码
import { Kafka } from "kafkajs"; import axios from "axios"; // Paketlerin nereden geldiğinin bilgisi const clientId = "my-app"; // Cluster Hangi Porta Yollayacak const brokers = ["kafka:9092"]; // Gonderilen Topic const topic = ["Topic2", "Topic3"]; const kafka = new Kafka({ clientId, brokers }); const consumer = kafka.consumer({ groupId: clientId }); let buffer = []; export const consume = async () => { await consumer.connect(); await consumer.subscribe({ topics: ["Topic2", "Topic3"] }); await consumer.run({ autoCommitInterval: 500, eachMessage: async ({ topic, partition, message }) => { const pkg = JSON.parse(message.value); if (buffer.length < 500) { buffer.push(pkg); } }, }); }; setInterval(async () => { const data = JSON.stringify(buffer); buffer = []; if (!data.length > 0) { console.log("İşlenecek veri olmadıgı için veri işleme alınmadı "); } else { try { await axios.post("http://db:3002/handle_data", data, { headers: { "Content-Type": "application/json", }, }); data = []; } catch (error) { console.error(`API error:${error}`); } } }, 200);
DB端代码
index.js文件运行在3002端口:
import { handle_data } from "./handle_data.js"; import express from "express"; //import sql from "mssql"; import { MongoClient } from "mongodb"; //const { MongoClient } = require("mongodb"); import cors from "cors"; const app = express(); app.use(cors()); app.use(express.json({ limit: "100mb" })); app.post("/handle_data", handle_data); app.listen(3002, () => { console.log("DB SİDE İS ALLİVE "); });
处理函数
import sql from "mssql"; export const handle_data = (req, res) => { const dataArray = req.body; if (dataArray.length > 0) { bulkInsert(req.body); } else { console.log( "Eklenecek Konum Verisi Olmadıgı İçin ekleme işlemi yapılamadı " ); } }; async function bulkInsert(data) { const config = { server: "A", database: "A", port: 1433, user: "A", password: "A", options: { trustServerCertificate: true, encrypt: false, }, }; try { await sql.connect(config); const table = new sql.Table("tblLocations"); table.create = false; table.columns.add("TaskId", sql.BigInt); table.columns.add("DoorNumber", sql.NVarChar(10)); table.columns.add("Latitude", sql.Float); table.columns.add("Longitude", sql.Float); table.columns.add("EmbeddedTime", sql.DateTime); table.columns.add("HatKodu", sql.NVarChar(10)); table.columns.add("GuzergahKodu", sql.NVarChar(50)); data.forEach((location) => { table.rows.add( location.taskId, location.doorNumber, location.Latitude, location.longitude, location.embeddedTime, location.hatKodu, location.routeCode ); }); const request = new sql.Request(); await request.bulk(table); console.log(` ${data.length} Adet Veri Veritabanına Eklendi ...`); } catch (err) { console.log(`Data uzunlugu ::::::::::::::::: ${data.length}`); console.error(err); } finally { } }
问题分析与解决方案
核心原因推测
- 端口耗尽(TIME_WAIT连接堆积):消费者每200ms发起一次HTTP请求,默认短连接会产生大量TIME_WAIT状态的连接,Linux系统可用端口有限,运行2小时后端口被占满,触发
EADDRNOTAVAIL。 - SQL连接泄漏:DB端的
bulkInsert函数未关闭SQL连接,导致连接池耗尽,无法处理新请求,间接引发消费者端连接失败。 - HTTP连接未正常释放:DB端处理函数未给客户端返回响应,导致HTTP连接长期挂起,加剧端口占用问题。
解决方案1:优化HTTP连接复用
修改消费者端Axios配置,开启Keep-Alive复用连接,减少TIME_WAIT连接数量:
import axios from "axios"; import http from "http"; // 创建带连接池的Axios实例 const axiosInstance = axios.create({ httpAgent: new http.Agent({ keepAlive: true, maxSockets: 10 }) // 限制并发连接数 }); // 替换原有axios.post为实例调用 await axiosInstance.post("http://db:3002/handle_data", data, { headers: { "Content-Type": "application/json", }, });
同时调整CentOS8内核参数,加快TIME_WAIT连接回收:
# 临时生效 echo 1 > /proc/sys/net/ipv4/tcp_tw_reuse echo 30 > /proc/sys/net/ipv4/tcp_fin_timeout echo 65535 > /proc/sys/net/ipv4/tcp_max_tw_buckets # 永久生效,编辑/etc/sysctl.conf添加以下内容 net.ipv4.tcp_tw_reuse = 1 net.ipv4.tcp_fin_timeout = 30 net.ipv4.tcp_max_tw_buckets = 65535 # 加载配置 sysctl -p
解决方案2:修复SQL连接泄漏与HTTP响应缺失
修改DB端bulkInsert函数,确保关闭SQL连接:
async function bulkInsert(data) { const config = { server: "A", database: "A", port: 1433, user: "A", password: "A", options: { trustServerCertificate: true, encrypt: false, }, }; let pool; try { pool = await sql.connect(config); const table = new sql.Table("tblLocations"); table.create = false; table.columns.add("TaskId", sql.BigInt); table.columns.add("DoorNumber", sql.NVarChar(10)); table.columns.add("Latitude", sql.Float); table.columns.add("Longitude", sql.Float); table.columns.add("EmbeddedTime", sql.DateTime); table.columns.add("HatKodu", sql.NVarChar(10)); table.columns.add("GuzergahKodu", sql.NVarChar(50)); data.forEach((location) => { table.rows.add( location.taskId, location.doorNumber, location.Latitude, location.longitude, location.embeddedTime, location.hatKodu, location.routeCode ); }); const request = new sql.Request(pool); await request.bulk(table); console.log(` ${data.length} Adet Veri Veritabanına Eklendi ...`); } catch (err) { console.log(`Data uzunlugu ::::::::::::::::: ${data.length}`); console.error(err); } finally { if (pool) await pool.close(); // 确保关闭连接池 } }
补充DB端处理函数的响应返回:
export const handle_data = (req, res) => { const dataArray = req.body; if (dataArray.length > 0) { bulkInsert(req.body) .then(() => res.status(200).send("数据插入成功")) .catch(err => res.status(500).send(`插入失败:${err.message}`)); } else { console.log("Eklenecek Konum Verisi Olmadıgı İçin ekleme işlemi yapılamadı "); res.status(200).send("无数据需要插入"); } };
解决方案3:修复消费者端代码逻辑错误
修改setInterval中的错误逻辑:
setInterval(async () => { const data = JSON.stringify(buffer); buffer = []; if (data.length === 0) { // 修正逻辑判断 console.log("İşlenecek veri olmadıgı için veri işleme alınmadı "); } else { try { await axiosInstance.post("http://db:3002/handle_data", data, { headers: { "Content-Type": "application/json", }, }); } catch (error) { console.error(`API error:${error}`); // 出错时将数据放回buffer,避免丢失 buffer.push(...JSON.parse(data)); } } }, 200);
内容的提问来源于stack exchange,提问作者Mehmet İMAL
相关产品推荐
相关产品推荐

