本文介绍如何在VPC环境下使用Node.js SDK接入消息队列Kafka版的默认接入点并收发消息。

前提条件

安装C++依赖库

  1. 执行以下命令切换到yum源配置目录/etc/yum.repos.d/
    cd /etc/yum.repos.d/
  2. 创建yum源配置文件confluent.repo
    [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. 执行以下命令安装C++依赖库。
    yum install librdkafka-devel

安装Node.js依赖库

  1. 执行以下命令为预处理器指定OpenSSL头文件路径。
    export CPPFLAGS=-I/usr/local/opt/openssl/include
  2. 执行以下命令为连接器指定OpenSSL库路径。
    export LDFLAGS=-L/usr/local/opt/openssl/lib
  3. 执行以下命令安装Node.js依赖库。
    npm install i --unsafe-perm node-rdkafka

准备配置

创建Kafka配置文件setting.js
module.exports = {
    'bootstrap_servers': ["kafka-ons-internet.aliyun.com:8080"],
    'topic_name': 'xxx',
    'consumer_id': 'xxx'
}
参数 描述
bootstrap_servers 默认接入点。您可在消息队列Kafka版控制台的实例详情页面的基本信息区域获取。
topic_name Topic名称。您可在消息队列Kafka版控制台的Topic管理页面获取。
consumer_id Consumer Group名称。您可在消息队列Kafka版控制台的Consumer Group管理页面获取。

发送消息

  1. 创建发送消息程序producer.js
    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
    });
    
    var connected = false
    
    producer.setPollInterval(100);
    
    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,      
          new Buffer('Hello Ali Kafka'),      
          null,   
          Date.now()
        );
      } 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");
  2. 执行以下命令发送消息。
    node producer.js

订阅消息

  1. 创建消费消息程序consumer.js
    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);
    });
  2. 执行以下命令消费消息。
    node consumer.js