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

node連接kafka2.0實現方法示例

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

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

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官網: https://kafka.apache.org/

相關文章

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

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

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

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

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

    詳解如何讓Express支持async/await

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

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

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

    更新Node.js的四種方法小結

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

    詳解用node.js實現簡單的反向代理

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

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

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

    Nodejs中怎么實現函數的串行執(zhí)行

    今天小編就為大家分享一篇關于Nodejs中怎么實現函數的串行執(zhí)行,小編覺得內容挺不錯的,現在分享給大家,具有很好的參考價值,需要的朋友一起跟隨小編來看看吧
    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實現熱部署等步驟
    2024-11-11
  • Node.js net模塊詳解(含類、方法、事件)

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

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

最新評論

临颍县| 亚东县| 平谷区| 保定市| 安仁县| 阆中市| 武安市| 峨边| 万宁市| 合水县| 竹溪县| 绥德县| 吉木萨尔县| 临朐县| 博乐市| 内丘县| 咸丰县| 漾濞| 西贡区| 绥阳县| 信丰县| 顺昌县| 宾阳县| 盐源县| 重庆市| 滨海县| 彭阳县| 阿瓦提县| 呼玛县| 兴仁县| 明水县| 宣汉县| 安塞县| 乌兰察布市| 屏东县| 望江县| 陆河县| 门源| 镇巴县| 德兴市| 普宁市|