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

Node.js TCP单条发送日志至Logstash仅首条生效问题求助

Node.js TCP单条发送日志到Logstash丢失问题解决

问题现象

  • 发送JSON数组格式日志时,Logstash能正常接收,但无法实现单条日志对应Elasticsearch索引的需求
  • 遍历数组逐个发送单条JSON日志时,仅第一条被Logstash接收,后续日志无迹可寻,且无任何错误输出
  • 尝试切换Logstash的json_lines codec,问题未解决

现有代码与配置

Node.js 发送代码

const client = new net.Socket();
client.connect(LOGSTASH_PORT, LOGSTASH_HOST, () => {

    for (const event of allEvents) {
      const success = client.write(JSON.stringify(event), (error) => {
        if (error) {
          console.error('Error writing to socket:', error);
        } else {
          console.log('Data written successfully');
        }
      });
    }

    //const success = client.write(JSON.stringify(allEvents));

    client.end();
});

Logstash 配置

input {
    tcp {
        port => "5400"
        codec =>"json"
    }
}

filter {

}

output {
 stdout { codec => json }
 elasticsearch {
   index => "<some index name>"
   hosts=> "${ELASTIC_HOSTS}"
   user=> "${ELASTIC_USER}"
   password=> "${ELASTIC_PASSWORD}"
   cacert=> "<some path>"
 }
}

问题根源

  1. TCP流无边界:TCP是字节流协议,Logstash的json codec会尝试将整个数据流解析为单个JSON对象。逐个发送的JSON没有分隔符,后续日志会被拼接成无效JSON,导致解析失败被丢弃
  2. Socket过早关闭:client.end()在所有异步write操作完成前就执行了,循环快速触发所有write后立刻关闭连接,后续数据还没发送就被中断

解决方法

1. 给单条日志添加换行符,配合json_lines codec

Logstash的json_lines codec以换行符作为单条日志的分隔符,发送每条JSON时末尾加上\n,让Logstash能正确识别每条日志:

修改Node.js代码:

const success = client.write(JSON.stringify(event) + '\n', (error) => {
  if (error) {
    console.error('Error writing to socket:', error);
  } else {
    console.log('Data written successfully');
  }
});

修改Logstash输入配置:

input {
    tcp {
        port => "5400"
        codec =>"json_lines"
    }
}

2. 等待所有写入完成后再关闭Socket

write是异步操作,不能循环后直接调用client.end(),必须确保所有数据发送完成再关闭连接,以下两种实现方式可选:

方式一:用Promise.all批量处理

const client = new net.Socket();
client.connect(LOGSTASH_PORT, LOGSTASH_HOST, async () => {
  // 把每个write操作包装成Promise
  const writeTasks = allEvents.map(event => {
    return new Promise((resolve, reject) => {
      client.write(JSON.stringify(event) + '\n', (error) => {
        error ? reject(error) : resolve();
        error ? console.error('Error writing to socket:', error) : console.log('Data written successfully');
      });
    });
  });

  // 等待所有写入完成
  await Promise.all(writeTasks);
  client.end();
});

方式二:通过计数跟踪完成状态

const client = new net.Socket();
client.connect(LOGSTASH_PORT, LOGSTASH_HOST, () => {
  let completedWrites = 0;
  const total = allEvents.length;

  for (const event of allEvents) {
    client.write(JSON.stringify(event) + '\n', (error) => {
      if (error) console.error('Error writing to socket:', error);
      else console.log('Data written successfully');
      
      completedWrites++;
      // 所有写入完成后关闭连接
      if (completedWrites === total) client.end();
    });
  }
});

3. 大日志量场景:处理背压避免丢失

如果日志量很大,client.write()可能返回false表示缓冲区已满,此时需要等待drain事件再继续写入:

const client = new net.Socket();
client.connect(LOGSTASH_PORT, LOGSTASH_HOST, () => {
  let currentIndex = 0;
  const totalEvents = allEvents.length;

  const sendNext = () => {
    if (currentIndex >= totalEvents) {
      client.end();
      return;
    }
    const event = allEvents[currentIndex];
    const isBufferAvailable = client.write(JSON.stringify(event) + '\n', (error) => {
      if (error) console.error('Error writing to socket:', error);
      else console.log('Data written successfully');
    });
    currentIndex++;
    // 缓冲区满,等待drain事件后继续发送
    if (!isBufferAvailable) {
      client.once('drain', sendNext);
    } else {
      sendNext();
    }
  };

  sendNext();
});

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 16:09:52