Node.js TCP单条发送日志至Logstash仅首条生效问题求助
Node.js TCP单条发送日志到Logstash丢失问题解决
问题现象
- 发送JSON数组格式日志时,Logstash能正常接收,但无法实现单条日志对应Elasticsearch索引的需求
- 遍历数组逐个发送单条JSON日志时,仅第一条被Logstash接收,后续日志无迹可寻,且无任何错误输出
- 尝试切换Logstash的
json_linescodec,问题未解决
现有代码与配置
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>" } }
问题根源
- TCP流无边界:TCP是字节流协议,Logstash的
jsoncodec会尝试将整个数据流解析为单个JSON对象。逐个发送的JSON没有分隔符,后续日志会被拼接成无效JSON,导致解析失败被丢弃 - 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
相关产品推荐
相关产品推荐

