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

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 {
  }
}

问题分析与解决方案

核心原因推测

  1. 端口耗尽(TIME_WAIT连接堆积):消费者每200ms发起一次HTTP请求,默认短连接会产生大量TIME_WAIT状态的连接,Linux系统可用端口有限,运行2小时后端口被占满,触发EADDRNOTAVAIL。
  2. SQL连接泄漏:DB端的bulkInsert函数未关闭SQL连接,导致连接池耗尽,无法处理新请求,间接引发消费者端连接失败。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 15:28:08