Use a biblioteca node-rdkafka para conectar uma aplicação Node.js ao ApsaraMQ for Kafka e produzir ou consumir mensagens.
Pré-requisitos
Antes de começar, verifique se você tem:
O GCC instalado
O Node.js 4.0.0 ou posterior instalado
O OpenSSL instalado
Uma instância do ApsaraMQ for Kafka com um tópico e um grupo de consumidores criados
(Apenas SSL) O certificado raiz SSL baixado e salvo como
ca-cert.pem
Instale a biblioteca C++
O pacote node-rdkafka depende da librdkafka, uma biblioteca nativa em C++. Instale-a pelo repositório Confluent.
-
Acesse o diretório de repositórios do yum:
cd /etc/yum.repos.d/ -
Crie um arquivo de configuração de repositório chamado
confluent.repocom o seguinte conteúdo:[Confluent.dist] name=Confluent repository (dist) baseurl=https://packages.confluent.io/rpm/5.1/7 gpgcheck=1 gpgkey=https://packages.confluent.io/rpm/5.1/archive.key enabled=1 [Confluent] name=Confluent repository baseurl=https://packages.confluent.io/rpm/5.1 gpgcheck=1 gpgkey=https://packages.confluent.io/rpm/5.1/archive.key enabled=1 -
Instale a
librdkafka:yum install librdkafka-devel
Instale a biblioteca Node.js
-
Defina os caminhos de cabeçalho e biblioteca do OpenSSL para a compilação nativa:
# Replace with your actual OpenSSL paths. # Run `whereis openssl` to find these paths. export CPPFLAGS=-I/usr/local/opt/openssl/include export LDFLAGS=-L/usr/local/opt/openssl/lib -
Instale o
node-rdkafka:npm install i --unsafe-perm node-rdkafka
Configure a conexão
Crie um arquivo setting.js com os parâmetros de conexão. Escolha a configuração correspondente ao tipo de endpoint.
Parâmetros de conexão
|
Parâmetro |
Descrição |
Onde encontrar |
|
|
Endpoint da instância do ApsaraMQ for Kafka |
Seção Endpoint Information na página Instance Details do console do ApsaraMQ for Kafka |
|
|
Nome do tópico para produção ou consumo de mensagens |
Página Topics no console do ApsaraMQ for Kafka |
|
|
ID do grupo de consumidores |
Página Groups no console do ApsaraMQ for Kafka |
|
|
Nome de usuário SASL (apenas para endpoint SSL) |
Campo Username na seção Configuration Information da página Instance Details |
|
|
Senha SASL (apenas para endpoint SSL) |
Campo Password na seção Configuration Information da página Instance Details |
Se a Access Control List (ACL) não estiver ativada na instância, obtenha o nome de usuário e a senha do Simple Authentication and Security Layer (SASL) nos campos Username e Password , na seção Configuration Information da página Instance Details no console do ApsaraMQ for Kafka . Caso a ACL esteja ativada, garanta que o usuário SASL tenha autorização para produzir e consumir mensagens. Para mais detalhes, consulte Conceder permissões a usuários SASL .
Endpoint padrão
module.exports = {
'bootstrap_servers': ["<bootstrap-servers>"],
'topic_name': '<topic-name>',
'consumer_id': '<consumer-group-id>'
}
Endpoint SSL
module.exports = {
'sasl_plain_username': '<sasl-username>',
'sasl_plain_password': '<sasl-password>',
'bootstrap_servers': ["<bootstrap-servers>"],
'topic_name': '<topic-name>',
'consumer_id': '<consumer-group-id>'
}
Após configurar o arquivo setting.js, envie-o junto com todos os arquivos da pasta demo correspondente para o servidor. Se usar o endpoint SSL, inclua o arquivo do certificado raiz SSL (ca-cert.pem) no mesmo diretório.
Dica: Baixe o projeto demo completo no repositório aliware-kafka-demos . Os exemplos em Node.js estão na pasta kafka-nodejs-demo , organizados por tipo de endpoint.
Produza mensagens
Salve o código a seguir como producer.js. Este exemplo envia uma mensagem a cada segundo e registra relatórios de entrega.
Execute o produtor:
node producer.js
Endpoint padrão
const Kafka = require('node-rdkafka');
const config = require('./setting');
console.log("features:" + Kafka.features);
console.log(Kafka.librdkafkaVersion);
var producer = new Kafka.Producer({
/*'debug': 'all', */
'api.version.request': 'true', // Negotiate API version with the broker automatically
'bootstrap.servers': config['bootstrap_servers'],
'dr_cb': true, // Enable delivery report callback
'dr_msg_cb': true // Include message payload in delivery reports
});
var connected = false
producer.setPollInterval(100); // Poll for delivery reports every 100 ms
producer.connect();
producer.on('ready', function() {
connected = true
console.log("connect ok")
});
producer.on("disconnected", function() {
connected = false;
producer.connect();
})
producer.on('event.log', function(event) {
console.log("event.log", event);
});
producer.on("error", function(error) {
console.log("error:" + error);
});
function produce() {
try {
producer.produce(
config['topic_name'],
null, // Partition (null = auto-assign)
new Buffer('Hello Ali Kafka'), // Message payload
null, // Message key
Date.now() // Timestamp
);
} catch (err) {
console.error('A problem occurred when sending our message');
console.error(err);
}
}
producer.on('delivery-report', function(err, report) {
console.log("delivery-report: producer ok");
});
producer.on('event.error', function(err) {
console.error('event.error:' + err);
})
setInterval(produce, 1000, "Interval");
Endpoint SSL
const Kafka = require('node-rdkafka');
const config = require('./setting');
console.log("features:" + Kafka.features);
console.log(Kafka.librdkafkaVersion);
var producer = new Kafka.Producer({
/*'debug': 'all', */
'api.version.request': 'true',
'bootstrap.servers': config['bootstrap_servers'],
'dr_cb': true,
'dr_msg_cb': true,
// SASL/SSL authentication
'security.protocol': 'sasl_ssl',
'ssl.ca.location': './ca-cert.pem', // Path to the SSL root certificate
'sasl.mechanisms': 'PLAIN',
'ssl.endpoint.identification.algorithm': 'none', // Disable hostname verification
'sasl.username': config['sasl_plain_username'],
'sasl.password': config['sasl_plain_password']
});
var connected = false
producer.setPollInterval(100);
producer.connect();
producer.on('ready', function() {
connected = true
console.log("connect ok")
});
function produce() {
try {
producer.produce(
config['topic_name'],
new Buffer('Hello Ali Kafka'),
null,
Date.now()
);
} catch (err) {
console.error('A problem occurred when sending our message');
console.error(err);
}
}
producer.on("disconnected", function() {
connected = false;
producer.connect();
})
producer.on('event.log', function(event) {
console.log("event.log", event);
});
producer.on("error", function(error) {
console.log("error:" + error);
});
producer.on('delivery-report', function(err, report) {
console.log("delivery-report: producer ok");
});
// Any errors we encounter, including connection errors
producer.on('event.error', function(err) {
console.error('event.error:' + err);
})
setInterval(produce, 1000, "Interval");
Consuma mensagens
Salve o código a seguir como consumer.js. Este exemplo assina um tópico e imprime cada mensagem no console.
Execute o consumidor:
node consumer.js
Endpoint padrão
const Kafka = require('node-rdkafka');
const config = require('./setting');
console.log(Kafka.features);
console.log(Kafka.librdkafkaVersion);
console.log(config)
var consumer = new Kafka.KafkaConsumer({
/*'debug': 'all',*/
'api.version.request': 'true',
'bootstrap.servers': config['bootstrap_servers'],
'group.id': config['consumer_id']
});
consumer.connect();
consumer.on('ready', function() {
console.log("connect ok");
consumer.subscribe([config['topic_name']]);
consumer.consume();
})
consumer.on('data', function(data) {
console.log(data);
});
consumer.on('event.log', function(event) {
console.log("event.log", event);
});
consumer.on('error', function(error) {
console.log("error:" + error);
});
consumer.on('event', function(event) {
console.log("event:" + event);
});
Endpoint SSL
const Kafka = require('node-rdkafka');
const config = require('./setting');
console.log(Kafka.features);
console.log(Kafka.librdkafkaVersion);
console.log(config)
var consumer = new Kafka.KafkaConsumer({
/*'debug': 'all',*/
'api.version.request': 'true',
'bootstrap.servers': config['bootstrap_servers'],
// SASL/SSL authentication
'security.protocol': 'sasl_ssl',
'ssl.endpoint.identification.algorithm': 'none',
'ssl.ca.location': './ca-cert.pem',
'sasl.mechanisms': 'PLAIN',
// Message size limits
'message.max.bytes': 32000,
'fetch.max.bytes': 32000,
'fetch.message.max.bytes': 32000,
'max.partition.fetch.bytes': 32000,
'sasl.username': config['sasl_plain_username'],
'sasl.password': config['sasl_plain_password'],
'group.id': config['consumer_id']
});
consumer.connect();
consumer.on('ready', function() {
console.log("connect ok");
consumer.subscribe([config['topic_name']]);
consumer.consume();
})
consumer.on('data', function(data) {
console.log(data);
});
consumer.on('event.log', function(event) {
console.log("event.log", event);
});
consumer.on('error', function(error) {
console.log("error:" + error);
});
consumer.on('event', function(event) {
console.log("event:" + event);
});