【发布时间】:2018-10-19 19:05:39
【问题描述】:
我有两台服务器,一台是 UI 服务器,另一台是 kafka 服务器。
我的 UI 服务器中的 javascript 文件正在从 csv 文件中获取数据,我正在逐行读取该文件并将其转换为 JSON。我需要将这些 JSON 格式的逐行读取数据发送到我的 Kafka 生产者服务器。为进一步工作。两台服务器都有自己的专用 IP 地址。例如kafka服务器有192.168.2.12:9098
reportJSON 是我在 UI 服务器 js 文件中获取 csv 数据的变量。
当我尝试运行 ui 服务器的 js 文件时,它显示错误:
2018-05-09T15:18:56.147Z - 错误:uncaughtException:io.connect 不是 函数日期=2018 年 5 月 9 日星期三 15:18:56 GMT+0000 (UTC)
JavaScript 文件内的 UI 连接:
var io = require('socket.io');
var socket = io.connect("http://192.168.2.12:9098");
socket.on('connect', function () {
console.log('Connection Established');
socket.emit('csvDataFromUI', function (reportJSON) {
console.log("Data inside the csvUpload Handler is = " + reportJSON);
});
});
kafka 生产者 javaScript 文件中的代码:
var http = require('http');
var app = express();
var host = process.env.HOST || config.host;
var port = process.env.PORT || config.port;
console.log("STARTING EVENT SERVER PRODUCER");
var server = http.createServer(app).listen(port, function () { });
server.timeout = 240000;
var io = require('socket.io').listen(server);
io.on('connection', function (socket) {
socket.on('csvDataFromUI', function(data) {
console.log("Data in kafka is = " + data);
});
//socket.emit('csvDataFromUI', payloadData);
});
/**************************************************** ******************************/ 新代码: 以下是:https://www.npmjs.com/package/kafka
UI 服务器 ui.js
将其创建为生产者:
var kafka = require('kafka');
var host = '192.168.2.12';
var port = 9098;
producer = new kafka.Producer({
host: host,
port: port,
topic: 'Postings',
partition: 0
});
producer.connect(function(reportJSON) {
console.log("rportJSON = " + reportJSON);
producer.send(reportJSON);
});
Kafka 服务器 kafkaProducer.js:
var kafkadata = require('kafka');
console.log("STARTING PRODUCER");
var consumer = new kafkadata.Consumer({
// these are the default values
host: '192.168.2.12',
port: 9098 ,
pollInterval: 2000,
maxSize: 1048576 // 1MB
})
consumer.on('message', function(topic, message) {
console.log(message)
})
consumer.connect(function() {
consumer.subscribeTopic({name: 'Postings', partition: 0})
})
我在 UI 服务器中遇到的错误是:error: uncaughtException: connect ECONNREFUSED reportJSON = 未定义
在 Kafka 服务器中,我看不到任何接收和获取错误: ReferenceError: 消息未定义
【问题讨论】:
标签: javascript node.js socket.io apache-kafka kafka-producer-api