Skip to content

Node.js Kafka

Kafka 是高吞吐量的分布式流处理平台,用于实时数据管道和流式应用,适合海量数据处理。

Kafka vs RabbitMQ

KafkaRabbitMQ
定位分布式流平台消息队列
吞吐量百万级/秒万级/秒
消息持久化默认持久化到磁盘需配置
消费模式拉取(pull)推送(push)
适合场景日志收集、大数据流处理异步任务、服务解耦

核心概念

概念说明
Producer生产者,写入消息
Consumer消费者,拉取消息
BrokerKafka 服务器节点
Topic消息分类主题
Partition主题分区(并行消费)
Consumer Group消费者组(组内分摊消费)

安装

bash
npm install kafkajs

生产者

js
import { Kafka } from 'kafkajs'

const kafka = new Kafka({ brokers: ['localhost:9092'] })
const producer = kafka.producer()

await producer.connect()
await producer.send({
  topic: 'order-topic',
  messages: [{ key: 'order-1', value: JSON.stringify({ id: 1, status: 'created' }) }]
})
await producer.disconnect()

消费者

js
const consumer = kafka.consumer({ groupId: 'order-group' })

await consumer.connect()
await consumer.subscribe({ topic: 'order-topic', fromBeginning: true })

await consumer.run({
  eachMessage: async ({ topic, partition, message }) => {
    console.log(`收到: ${message.value.toString()}`)
  }
})

典型场景

场景说明
日志收集各服务日志统一汇入 Kafka
实时流处理点击流、交易流实时分析
事件溯源业务事件持久化,可回放
数据管道数据库 CDC 同步到大屏/数仓