Compartir tecnología

Nodejs Capítulo 80 (Kafka Avanzado)

2024-07-12

한어Русский языкEnglishFrançaisIndonesianSanskrit日本語DeutschPortuguêsΕλληνικάespañolItalianoSuomalainenLatina

Insertar descripción de la imagen aquí

KafkaLos conocimientos previos se han analizado en capítulos anteriores y no se repetirán nuevamente.

Operaciones del clúster Kafka

1. Cree múltiples servicios Kafka

Hacer una copiakafkaDirectorio completo renombradokafka2

RevisarArchivo de configuración kafka2/config/server.properties Este archivo

broker.id=1 //唯一broker
port=9093 //切换端口
listeners=PLAINTEXT://:9093 //切换监听源
  • 1
  • 2
  • 3

puesta en marchaGuardián del zoológicoy kafka y kafka2

.binwindowskafka-server-start.bat .configserver.properties
  • 1

2.Gestión de clientes

Ver información del clúster y objetos del cliente.

import { Kafka, CompressionTypes } from 'kafkajs'

const kafka = new Kafka({
    clientId: 'my-app', //客户端标识
    brokers: ['localhost:9092', 'localhost:9093'], //kafka集群
})

const admin = kafka.admin() //创建admin对象
await admin.connect() //连接kafka
const cluster = await admin.describeCluster() //获取集群信息
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10

El valor de retorno se puede utilizar para ver información sobre la conexión al clúster, como el ID del puerto, etc.

{
  brokers: [
    { nodeId: 0, host: '26.26.26.1', port: 9092 },
    { nodeId: 1, host: '26.26.26.1', port: 9093 }
  ],
  controller: 0,
  clusterId: 'XHa77me4TZWO8cfWSTHoaQ'
}
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8

Crear temacreateTopicsanalizarátrueSi el tema se creó correctamente ofalse ¿Ya existe?Si ocurre un error, este método generará una excepción.

Eliminar temaadmin.deleteTopics Pasar temas eliminados

Ver lista de temaslistTopics Enumera los nombres de todos los temas existentes y devuelve una serie de cadenas. Si ocurre un error, este método generará una excepción.

//创建主题
await admin.createTopics({
    topics: [
        { topic: 'xiaoman', numPartitions: 1, replicationFactor: 1 },
        { topic: 'xiaoman2', numPartitions: 1, replicationFactor: 1 },
    ],
})
//删除主题
await admin.deleteTopics({ topics: ['xiaoman', 'xiaoman2'] })
//查看主题
await admin.listTopics().then(topics => {
    console.log('topics', topics)
})
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • 11
  • 12
  • 13
3. Asuntos

KafkaJS brinda soporte para transacciones Kafka, que se pueden utilizar para realizar operaciones con características transaccionales. Las transacciones Kafka se utilizan para garantizar que un grupo de mensajes relacionados sea全部成功提交要么全部回滚, manteniendo así la coherencia de los datos

import { Kafka, CompressionTypes } from 'kafkajs'

const kafka = new Kafka({
    clientId: 'my-app', //客户端标识
    brokers: ['localhost:9092', 'localhost:9093'], //kafka集群
})

//生产者
const producer = kafka.producer({
    transactionalId: '填写事务ID',
    maxInFlightRequests: 1, //最大同时发送请求数
    idempotent: true, //是否开启幂等提交
})
//连接服务器
await producer.connect()

const transaction = await producer.transaction()
try {
    await transaction.send({
        topic: 'xiaoman',
        messages: [{ value: '100元' }],
    })
    await transaction.commit() // 事务提交
}
catch (e) {
    console.log(e)
    await transaction.abort() // 事务提交失败,回滚
}
await admin.disconnect()
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • 11
  • 12
  • 13
  • 14
  • 15
  • 16
  • 17
  • 18
  • 19
  • 20
  • 21
  • 22
  • 23
  • 24
  • 25
  • 26
  • 27
  • 28
  • 29