Node.js Kafka
Kafka 是高吞吐量的分布式流处理平台,用于实时数据管道和流式应用,适合海量数据处理。
Kafka vs RabbitMQ
| Kafka | RabbitMQ | |
|---|---|---|
| 定位 | 分布式流平台 | 消息队列 |
| 吞吐量 | 百万级/秒 | 万级/秒 |
| 消息持久化 | 默认持久化到磁盘 | 需配置 |
| 消费模式 | 拉取(pull) | 推送(push) |
| 适合场景 | 日志收集、大数据流处理 | 异步任务、服务解耦 |
核心概念
| 概念 | 说明 |
|---|---|
| Producer | 生产者,写入消息 |
| Consumer | 消费者,拉取消息 |
| Broker | Kafka 服务器节点 |
| 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 同步到大屏/数仓 |