All Products
Search
Document Center

ApsaraMQ for Kafka:Gunakan SDK untuk Node.js untuk mengirim dan menerima pesan

Last Updated:Jul 06, 2025

Topik ini menjelaskan cara menggunakan SDK untuk Node.js untuk terhubung ke ApsaraMQ for Kafka guna mengirim dan menerima pesan.

Persyaratan lingkungan

  • GNU Compiler Collection (GCC) telah diinstal. Untuk informasi lebih lanjut, lihat Menginstal GCC.

  • Node.js telah diinstal. Untuk informasi lebih lanjut, lihat Unduhan.

    Penting

    Versi Node.js harus 4.0.0 atau yang lebih baru.

  • OpenSSL telah diinstal. Untuk informasi lebih lanjut, lihat Unduhan.

Instal pustaka C++

  1. Jalankan perintah berikut untuk beralih ke direktori repositori yum /etc/yum.repos.d/:

    cd /etc/yum.repos.d/
  2. Buat file konfigurasi repositori yum bernama 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. Jalankan perintah berikut untuk menginstal pustaka C++:

    yum install librdkafka-devel

Instal pustaka Node.js

  1. Jalankan perintah berikut untuk menentukan jalur file header OpenSSL untuk preprocessor:

    # Setel parameter ini ke jalur file header OpenSSL yang diinstal pada mesin lokal Anda.
    export CPPFLAGS=-I</usr/local/opt/openssl/include>
  2. Jalankan perintah berikut untuk menentukan jalur pustaka OpenSSL untuk konektor:

    # Setel parameter ini ke jalur file pustaka OpenSSL yang diinstal pada mesin lokal Anda.
    export LDFLAGS=-L</usr/local/opt/openssl/lib>
  3. Jalankan perintah berikut untuk menginstal pustaka Node.js:

    npm install i --unsafe-perm node-rdkafka
Catatan

Jalankan perintah whereis openssl untuk mendapatkan jalur file header OpenSSL dan jalur pustaka OpenSSL.

Buat file konfigurasi

  1. (Opsional) Unduh sertifikat root Secure Sockets Layer (SSL). Jika Anda menggunakan titik akhir SSL untuk terhubung ke instance ApsaraMQ for Kafka, Anda harus menginstal sertifikat.

  2. Pergi ke halaman aliware-kafka-demos, klik download untuk mengunduh proyek demo ke mesin lokal Anda, lalu ekstrak paket proyek demo tersebut.

  3. Di dalam paket yang diekstraksi, pergi ke folder kafka-nodejs-demo. Kemudian, buka folder yang sesuai berdasarkan titik akhir yang ingin digunakan dan konfigurasikan file setting.js di dalam folder tersebut.

    module.exports = {
        'sasl_plain_username': 'XXX',
        'sasl_plain_password': 'XXX',
        'bootstrap_servers': ["XXX"],
        'topic_name': 'XXX',
        'consumer_id': 'XXX'
    }

    Parameter

    Deskripsi

    sasl_plain_username

    Nama pengguna Simple Authentication and Security Layer (SASL). Jika Anda menggunakan titik akhir default untuk terhubung ke instance ApsaraMQ for Kafka, parameter ini tidak tersedia.

    Catatan
    • Jika fitur ACL tidak diaktifkan untuk instance ApsaraMQ for Kafka, Anda dapat memperoleh nama pengguna dan kata sandi pengguna SASL dari parameter Username dan Password di bagian Configuration Information pada halaman Instance Details di Konsol ApsaraMQ for Kafka.

    • Jika fitur ACL diaktifkan untuk instance ApsaraMQ for Kafka, pastikan bahwa pengguna SASL diberi otorisasi untuk mengirim dan menerima pesan menggunakan instance tersebut. Untuk informasi lebih lanjut, lihat Berikan izin kepada pengguna SASL.

    sasl_plain_password

    Kata sandi pengguna SASL. Jika Anda menggunakan titik akhir default untuk terhubung ke instance ApsaraMQ for Kafka, parameter ini tidak tersedia.

    bootstrap_servers

    Titik akhir SSL dari instance ApsaraMQ for Kafka. Anda dapat memperoleh titik akhir di bagian Endpoint Information pada halaman Instance Details di Konsol ApsaraMQ for Kafka.

    topic_name

    Nama topik. Anda dapat memperoleh nama topik di halaman Topics di Konsol ApsaraMQ for Kafka.

    consumer_id

    ID grup. Anda dapat memperoleh ID grup di halaman Groups di Konsol ApsaraMQ for Kafka.

  4. Setelah parameter yang diperlukan dikonfigurasi, unggah semua file di folder tempat file konfigurasi berada ke direktori instalasi pustaka dependensi Node.js di server Anda. Folder yang sesuai dengan titik akhir SSL berisi file sertifikat root SSL.

Kirim pesan

Jalankan perintah berikut untuk menjalankan producer.js guna mengirim pesan:

node producer.js

Kode sampel berikut memberikan contoh consumer.js:

  • Jika Anda menggunakan titik akhir default untuk terhubung ke instance ApsaraMQ for Kafka, gunakan kode sampel berikut:

    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");
  • Jika Anda menggunakan titik akhir SSL untuk terhubung ke instance ApsaraMQ for Kafka, gunakan kode sampel berikut:

    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,
        'security.protocol' : 'sasl_ssl',
        'ssl.ca.location' : './ca-cert.pem',
        'sasl.mechanisms' : 'PLAIN',
        'ssl.endpoint.identification.algorithm':'none',
        '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");

Terima pesan

Jalankan perintah berikut untuk menjalankan consumer.js guna menerima pesan:

node consumer.js

Kode sampel berikut memberikan contoh consumer.js:

  • Jika Anda menggunakan titik akhir default untuk terhubung ke instance ApsaraMQ for Kafka, gunakan kode sampel berikut:

    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);
    });
  • Jika Anda menggunakan titik akhir SSL untuk terhubung ke instance ApsaraMQ for Kafka, gunakan kode sampel berikut:

    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'],
        'security.protocol' : 'sasl_ssl',
        'ssl.endpoint.identification.algorithm':'none',
        'ssl.ca.location' : './ca-cert.pem',
        'sasl.mechanisms' : 'PLAIN',
        '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);
    });