Todos os produtos
Search
Central de documentação

ApsaraMQ for Kafka:Send and receive messages with the Node.js SDK

Última atualização: Jun 27, 2026

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.

  1. Acesse o diretório de repositórios do yum:

       cd /etc/yum.repos.d/
  2. Crie um arquivo de configuração de repositório chamado confluent.repo com 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
  3. Instale a librdkafka:

       yum install librdkafka-devel

Instale a biblioteca Node.js

  1. 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
  2. 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

bootstrap_servers

Endpoint da instância do ApsaraMQ for Kafka

Seção Endpoint Information na página Instance Details do console do ApsaraMQ for Kafka

topic_name

Nome do tópico para produção ou consumo de mensagens

Página Topics no console do ApsaraMQ for Kafka

consumer_id

ID do grupo de consumidores

Página Groups no console do ApsaraMQ for Kafka

sasl_plain_username

Nome de usuário SASL (apenas para endpoint SSL)

Campo Username na seção Configuration Information da página Instance Details

sasl_plain_password

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);
});