终于解决了。我犯了一个错误。默认情况下 Apache Kafka 创建 3 个复制,这意味着默认情况下它创建 3 个代理。然而,上面 YAML 只创建一个代理,并寻找其他 2 个未创建的代理。因此我们得到了同样的错误。
修复
version: '3'
services:
zookeeper:
container_name: zookeeper
image: confluentinc/cp-zookeeper
ports:
- "32181:32181"
environment:
ZOOKEEPER_CLIENT_PORT: 32181
ZOOKEEPER_TICK_TIME: 2000
ZOOKEEPER_SYNC_LIMIT: 2
kafka:
container_name: kafka
image: confluentinc/cp-kafka
ports:
- "9094:9094"
environment:
KAFKA_ZOOKEEPER_CONNECT: zookeeper:32181
KAFKA_LISTENERS: INTERNAL://:9092,OUTSIDE://:9094
KAFKA_ADVERTISED_LISTENERS: INTERNAL://:9092,OUTSIDE://localhost:9094
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,OUTSIDE:PLAINTEXT
KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
ES_JAVA_OPTS: "-Xms512m -Xmx3000m"
然后KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 告诉 Apache kafka 我们只想创建一个复制及其工作。
使用以下代码发布和消费kafka消息
var intCounter = 104;
setInterval(() => {
setMessageInTopic('topic'+intCounter, function(response){
echo("Producer Response"+response);
if(response){
echo("Consumer init");
getMessageFromKafak(response,function(pStrMessage){
echo ("------------------------Consumer message --------------------");
echo(pStrMessage);
});
}else{
echo("Error");
}
});
intCounter++;
},5000);
function setMessageInTopic(pStrTopicName, callback){
const kafka = require('kafka-node');
try {
const Producer = kafka.Producer;
const client = new kafka.KafkaClient({kafkaHost:"localhost:9094"});
const producer = new Producer(client);
const kafka_topic = pStrTopicName;
console.log(kafka_topic);
let payloads = [
{
topic: kafka_topic,
messages: {'name':'vipin'}
}
];
producer.on('ready', async function() {
let push_status = producer.send(payloads, (err, data) => {
if (err) {
console.log('[kafka-producer -> '+kafka_topic+']: broker update failed');
callback(false);
} else {
console.log('[kafka-producer -> '+kafka_topic+']: broker update success');
callback(pStrTopicName);
}
});
});
producer.on('error', function(err) {
console.log(err);
console.log('[kafka-producer -> '+kafka_topic+']: connection errored');
callback(false);
///throw err;
});
}
catch(e) {
console.log(e);
callback(false);
}
}
function getMessageFromKafak(pStrTopicName, callback){
try{
var kafka = require('kafka-node');
//var HighLevelProducer = kafka.HighLevelProducer;
var Consumer = kafka.Consumer;
var client = new kafka.KafkaClient({kafkaHost:"localhost:9094"});
let consumer = new Consumer(
client,
[{ topic: pStrTopicName}],
{
autoCommit: true,
fetchMaxWaitMs: 1000,
fetchMaxBytes: 1024 * 1024,
encoding: 'utf8',
fromOffset: false
}
);
echo("------------------RESPONE MESSAGE PROCESS "+pStrTopicName+" ----------------------")
consumer.on('message', function (message) {
echo("------------------RESPONE MESSAGE ----------------------")
console.log(message);
consumer.close();
client.close();
return callback(message);
}).on('error', function (message) {
echo("------------------RESPONE MESSAGE ERROR ----------------------")
console.log(message);
consumer.close();
client.close();
return callback(message);
});
}catch(e) {
console.log(e);
}
}
编码愉快:)
感谢和问候
Jaiswar Vipin Kumar R.