最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

node連接kafka2.0實現(xiàn)方法示例

 更新時間:2023年05月26日 09:14:18   作者:他強任他強03  
這篇文章主要介紹了node連接kafka2.0,nodejs連接kafka2.0的實現(xiàn)方法,結(jié)合實例形式分析了kafka2.0的功能、原理、以及node.js連接kafka2.0的具體實現(xiàn)技巧,需要的朋友可以參考下

Kafka是由Apache軟件基金會開發(fā)的一個開源流處理平臺,由Scala和Java編寫。Kafka是一種高吞吐量的分布式發(fā)布訂閱消息系統(tǒng),它可以處理消費者在網(wǎng)站中的所有動作流數(shù)據(jù)

node.js使用Kafka需要安裝的npm包:https://www.npmjs.com/package/wisrtoni40-confluent-schema#Quickstart

npm i wisrtoni40-confluent-schema --save

procedurer.ts文件

import { HighLevelProducer, KafkaClient } from 'kafka-node';
 import { v4 as uuidv4 } from 'uuid';
 import {
   ConfluentAvroStrategy,
   ConfluentMultiRegistry,
   ConfluentPubResolveStrategy,
 } from 'wisrtoni40-confluent-schema';
 /**
  * -----------------------------------------------------------------------------
  * Config
  * -----------------------------------------------------------------------------
  */
 const kafkaHost = '你的kafka host';
 const topic = '你的topic';
 const registryHost =
   '你的kafka注冊host';
 /**
  * -----------------------------------------------------------------------------
  * Kafka Client and Producer
  * -----------------------------------------------------------------------------
  */
 const kafkaClient = new KafkaClient({
   kafkaHost,
   clientId: uuidv4(),
   connectTimeout: 60000,
   requestTimeout: 60000,
   connectRetryOptions: {
     retries: 5,
     factor: 0,
     minTimeout: 1000,
     maxTimeout: 1000,
     randomize: false,
   },
   sasl: {
     mechanism: 'plain',
     username: '你的kafka用戶名',
     password: '你的kafka密碼',
   },
 });
 const producer = new HighLevelProducer(kafkaClient, {
   requireAcks: 1,
   ackTimeoutMs: 100,
 });
 /**
  * -----------------------------------------------------------------------------
  * Confluent Resolver
  * -----------------------------------------------------------------------------
  */
 const schemaRegistry = new ConfluentMultiRegistry(registryHost);
 const avro = new ConfluentAvroStrategy();
 const resolver = new ConfluentPubResolveStrategy(schemaRegistry, avro, topic);
 /**
  * -----------------------------------------------------------------------------
  * Produce
  * -----------------------------------------------------------------------------
  */
 (async () => {
   const data = {
    evt_dt: 1664446229425,
    evt_type: 'tower_unload',
    plant: 'F110',
    machineName: 'TOWER_01',
    errorCode: '',
    description: '',
    result: 'OK',
    evt_ns: 'wmy.dx',
    evt_tp: 'tower.error',
    evt_pid: 'TOWER_01',
    evt_pubBy: 'nifi.11142'
   };
   const processedData = await resolver.resolve(data);
   producer.send([{ topic, messages: processedData }], (error, result) => {
     if (error) {
       console.error(error);
     } else {
       console.log(result);
     }
   });
 })();

procedurer.js文件

var kafka_node_1 = require("kafka-node");
var uuid_1 = require("uuid");
var wisrtoni40_confluent_schema_1 = require("wisrtoni40-confluent-schema");
var kafkaHost = '你的kafka host';
var topic = '你的topic';
var registryHost = '你的kafka注冊host';
const kafkaClient = new kafka_node_1.KafkaClient({
  kafkaHost,
  clientId: (0, uuid_1.v4)(),
  connectTimeout: 60000,
  requestTimeout: 60000,
  connectRetryOptions: {
    retries: 5,
    factor: 0,
    minTimeout: 1000,
    maxTimeout: 1000,
    randomize: false,
  },
  sasl: {
    mechanism: 'plain',
    username: '你的kafka用戶名',
    password: '你的kafka密碼',
  },
});
const producer = new kafka_node_1.HighLevelProducer(kafkaClient, {
  requireAcks: 1,
  ackTimeoutMs: 100,
});
const schemaRegistry = new wisrtoni40_confluent_schema_1.ConfluentMultiRegistry(registryHost);
const avro = new wisrtoni40_confluent_schema_1.ConfluentAvroStrategy();
const resolver = new wisrtoni40_confluent_schema_1.ConfluentPubResolveStrategy(schemaRegistry, avro, topic);
(async () => {
  const data = {
    evt_dt: 1664446229425,
    evt_type: 'tower_unload',
    plant: 'F110',
    machineName: 'TOWER_01',
    errorCode: '',
    description: '',
    result: 'OK',
    evt_ns: 'wmy.dx',
    evt_tp: 'tower.error',
    evt_pid: 'TOWER_01',
    evt_pubBy: 'nifi.11142'
  };
  const processedData = await resolver.resolve(data);
  producer.send([{ topic, messages: processedData }], (error, result) => {
    if (error) {
      console.error(error);
    } else {
      console.log(result);
    }
  });
})();

consumer.ts文件

import { ConsumerGroup } from 'kafka-node';
import { v4 as uuidv4 } from 'uuid';
import {
  ConfluentAvroStrategy,
  ConfluentMultiRegistry,
  ConfluentSubResolveStrategy,
} from 'wisrtoni40-confluent-schema';
/**
 * -----------------------------------------------------------------------------
 * Config
 * -----------------------------------------------------------------------------
 */
const kafkaHost = '你的kafka host';
const topic = '你的topic';
const registryHost =
  '你的kafka注冊host';
/**
 * -----------------------------------------------------------------------------
 * Kafka Consumer
 * -----------------------------------------------------------------------------
 */
const consumer = new ConsumerGroup(
  {
    kafkaHost,
    groupId: uuidv4(),
    sessionTimeout: 15000,
    protocol: ['roundrobin'],
    encoding: 'buffer',
    fromOffset: 'latest',
    outOfRangeOffset: 'latest',
    sasl: {
      mechanism: 'plain',
      username: '你的kafka用戶名',
      password: '你的kafka密碼',
    },
  },
  topic,
);
/**
 * -----------------------------------------------------------------------------
 * Confluent Resolver
 * -----------------------------------------------------------------------------
 */
const schemaRegistry = new ConfluentMultiRegistry(registryHost);
const avro = new ConfluentAvroStrategy();
const resolver = new ConfluentSubResolveStrategy(schemaRegistry, avro);
/**
 * -----------------------------------------------------------------------------
 * Consume
 * -----------------------------------------------------------------------------
 */
consumer.on('message', async msg => {
  const result = await resolver.resolve(msg.value);
  console.log(msg.offset);
  console.log(result);
});

comsumer.js文件

var kafka_node_1 = require("kafka-node");
var uuid_1 = require("uuid");
var wisrtoni40_confluent_schema_1 = require("wisrtoni40-confluent-schema");
var kafkaHost = '你的kafka host';
var topic = '你的topic';
var registryHost = '你的kafka注冊host';
var consumer = new kafka_node_1.ConsumerGroup({
  kafkaHost: kafkaHost,
  groupId: (0, uuid_1.v4)(),
  sessionTimeout: 15000,
  protocol: ['roundrobin'],
  encoding: 'buffer',
  fromOffset: 'latest',
  outOfRangeOffset: 'latest',
  sasl: {
    mechanism: 'plain',
    username: '你的kafka用戶名',
    password: '你的kafka密碼'
  }
}, topic);
var schemaRegistry = new wisrtoni40_confluent_schema_1.ConfluentMultiRegistry(registryHost);
var avro = new wisrtoni40_confluent_schema_1.ConfluentAvroStrategy();
var resolver = new wisrtoni40_confluent_schema_1.ConfluentSubResolveStrategy(schemaRegistry, avro);
consumer.on('message', async function (msg) {
  const result = await resolver.resolve(msg.value);
  console.log(msg.offset);
  console.log(result);
});

附:kafka官網(wǎng): https://kafka.apache.org/

相關文章

  • 深入解讀Node.js中的koa源碼

    深入解讀Node.js中的koa源碼

    這篇文章主要介紹了深入解讀Node.js中的koa源碼,任何一個框架的出現(xiàn)都是為了解決問題,而Koa則是為了更方便的構(gòu)建http服務而出現(xiàn)的。 可以簡單的理解為一個HTTP服務的中間件框架。,需要的朋友可以參考下
    2019-06-06
  • Koa2微信公眾號開發(fā)之本地開發(fā)調(diào)試環(huán)境搭建

    Koa2微信公眾號開發(fā)之本地開發(fā)調(diào)試環(huán)境搭建

    本篇文章主要介紹了Koa2微信公眾號開發(fā)之本地開發(fā)調(diào)試環(huán)境搭建,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2018-05-05
  • 詳解如何讓Express支持async/await

    詳解如何讓Express支持async/await

    本篇文章主要介紹了詳解如何讓Express支持async/await,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-10-10
  • NodeJS實現(xiàn)圖片上傳代碼(Express)

    NodeJS實現(xiàn)圖片上傳代碼(Express)

    本篇文章主要介紹了NodeJS實現(xiàn)圖片上傳代碼(Express) ,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-06-06
  • 更新Node.js的四種方法小結(jié)

    更新Node.js的四種方法小結(jié)

    Node.js是一個開放源代碼的跨平臺JavaScript運行環(huán)境,它在不同的平臺上都得到了廣泛使用和支持,強大的生態(tài)系統(tǒng)、持續(xù)的更新和不斷改進的性能使得Node.js非常受歡迎,然而,更新Node.js仍然是一個必要的過程,本文給大家介紹一些有關如何更新Node.js的方法
    2023-11-11
  • 詳解用node.js實現(xiàn)簡單的反向代理

    詳解用node.js實現(xiàn)簡單的反向代理

    本篇文章主要介紹了詳解用node.js實現(xiàn)簡單的反向代理,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-06-06
  • node.js?express和koa中間件機制和錯誤處理機制

    node.js?express和koa中間件機制和錯誤處理機制

    這篇文章主要介紹了node.js?express和koa中間件機制和錯誤處理機制,文章圍繞主題展開詳細的內(nèi)容介紹,具有一定的參考價值,需要的朋友可以參考一下
    2022-07-07
  • Nodejs中怎么實現(xiàn)函數(shù)的串行執(zhí)行

    Nodejs中怎么實現(xiàn)函數(shù)的串行執(zhí)行

    今天小編就為大家分享一篇關于Nodejs中怎么實現(xiàn)函數(shù)的串行執(zhí)行,小編覺得內(nèi)容挺不錯的,現(xiàn)在分享給大家,具有很好的參考價值,需要的朋友一起跟隨小編來看看吧
    2019-03-03
  • 教你從零開始在Windows系統(tǒng)上搭建一個node.js后端服務項目

    教你從零開始在Windows系統(tǒng)上搭建一個node.js后端服務項目

    這篇文章詳細介紹了如何在Windows環(huán)境下搭建一個Node.js項目并使用Express框架,包括安裝Node.js、配置環(huán)境、創(chuàng)建項目、安裝Express、編輯代碼、運行項目、集成Nodemon實現(xiàn)熱部署等步驟
    2024-11-11
  • Node.js net模塊詳解(含類、方法、事件)

    Node.js net模塊詳解(含類、方法、事件)

    Node.js 的 net 模塊提供了基于 TCP 或 IPC 的網(wǎng)絡通信能力,用于創(chuàng)建服務器和客戶端,本文給大家介紹Node.js net模塊詳解包含類、方法、事件及示例,感興趣的朋友一起看看吧
    2025-04-04

最新評論

奉节县| 楚雄市| 苏尼特左旗| 瓦房店市| 临沭县| 绥中县| 和顺县| 田东县| 青岛市| 林周县| 包头市| 太康县| 永和县| 平乐县| 威宁| 湖州市| 四川省| 崇信县| 无棣县| 喀喇沁旗| 津市市| 公主岭市| 成都市| 壶关县| 南通市| 宁海县| 赤城县| 成武县| 宁武县| 成都市| 吴忠市| 咸宁市| 微山县| 漠河县| 个旧市| 嵊泗县| 康乐县| 东源县| 和林格尔县| 高密市| 彭阳县|