两台服务器间Socket.io连接与Kafka数据传输故障排查
解决UI服务器与Kafka服务器数据传输的两个阶段问题
一、Socket.IO连接阶段:io.connect is not a function错误
错误原因
你在UI服务器代码里误用了服务端的socket.io包做客户端连接,但客户端必须使用socket.io-client包——服务端的socket.io并不提供io.connect()方法,这直接导致了报错。另外你的emit调用逻辑也有问题,第二个参数应该直接传reportJSON,而非包裹成函数。
修正步骤与代码
- 先在UI服务器安装客户端依赖:
npm install socket.io-client
- 修改UI服务器的连接代码:
var io = require('socket.io-client'); // 直接传入服务端地址创建连接 var socket = io("http://192.168.2.12:9098"); socket.on('connect', function () { console.log('Connection Established'); // 直接将reportJSON作为emit的第二个参数发送 socket.emit('csvDataFromUI', reportJSON, function () { console.log("Data sent to Kafka server: " + reportJSON); }); });
- 确认Kafka服务器端依赖完整(先安装所需包):
npm install express socket.io
二、使用kafka npm包阶段的两个错误
错误1:UI服务器connect ECONNREFUSED + reportJSON = undefined
原因
ECONNREFUSED:说明UI服务器无法连接到Kafka服务器的9098端口,大概率是Kafka默认监听端口为9092而非9098、服务器防火墙拦截了端口,或者你用的kafka包过于老旧,对新版本Kafka兼容性差。reportJSON = undefined:producer.connect()的回调函数没有reportJSON这个参数,它只是连接成功的通知,你需要在外部先拿到转换好的reportJSON,再在连接成功后发送。
修正方案
推荐用维护活跃、兼容性更好的kafkajs替代老旧的kafka包,同时检查Kafka端口配置:
- 安装
kafkajs:
npm install kafkajs
- 修改UI服务器的生产者代码(ui.js):
const { Kafka } = require('kafkajs') const kafka = new Kafka({ clientId: 'ui-producer', brokers: ['192.168.2.12:9092'] // 若Kafka确实改了9098端口,请对应修改 }) const producer = kafka.producer() // 确保已完成CSV转JSON逻辑,拿到reportJSON后执行发送 async function sendData() { try { await producer.connect() await producer.send({ topic: 'Postings', messages: [ { value: JSON.stringify(reportJSON) } // Kafka消息需为字符串格式 ], }) console.log("Data sent successfully: ", reportJSON) } catch (err) { console.error("Failed to send data: ", err) } finally { await producer.disconnect() } } sendData()
错误2:Kafka服务器ReferenceError: message is not defined
原因
你使用的kafka包中,consumer.on('message')的回调参数格式错误——正确格式是单个message对象,而非(topic, message),这导致第二个参数message未被定义,触发报错。
修正方案(同步使用kafkajs)
修改Kafka服务器的消费者代码(kafkaProducer.js):
const { Kafka } = require('kafkajs') const kafka = new Kafka({ clientId: 'kafka-consumer', brokers: ['192.168.2.12:9092'] // 对应Kafka实际端口 }) const consumer = kafka.consumer({ groupId: 'postings-group' }) async function run() { console.log("STARTING CONSUMER") await consumer.connect() await consumer.subscribe({ topic: 'Postings', fromBeginning: true }) await consumer.run({ eachMessage: async ({ topic, partition, message }) => { console.log("Received data: ", JSON.parse(message.value.toString())) }, }) } run().catch(console.error)
内容的提问来源于stack exchange,提问作者Shalabh Negi
相关产品推荐
相关产品推荐

