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

两台服务器间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,而非包裹成函数。

修正步骤与代码

  1. 先在UI服务器安装客户端依赖:
npm install socket.io-client
  1. 修改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);
  });
});
  1. 确认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端口配置:

  1. 安装kafkajs:
npm install kafkajs
  1. 修改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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 06:43:11