事件驱动架构:Kafka 开发指南
深入讲解 Kafka 在事件驱动架构中的应用,包含实际代码示例。经典后端技术,与 AI 无关。
深入讲解 Kafka 在事件驱动架构中的应用,包含实际代码示例。经典后端技术,与 AI 无关。
你好,我是 Maneshwar。我正在开发 git-lrc,一个在每次提交时运行的 AI 代码审查工具。它是免费的,代码开源在 Github 上。请给我们 Star,帮助开发者发现这个项目。也欢迎你试用并分享反馈,帮助我们改进产品。
Kafka 最初在 LinkedIn 开发,后来在 Apache 软件基金会开源,现已广泛用于构建高吞吐量、容错、可扩展的数据管道、实时分析和事件驱动架构。
在 Kafka 出现之前,RabbitMQ 和 ActiveMQ 等传统消息队列被广泛使用,但它们在处理大规模、高吞吐量的实时数据流时存在局限。
Kafka 被设计用来解决这些问题,提供以下功能:
大规模数据处理 – Kafka 优化用于在分布式系统中摄取、存储和分发高容量的数据流。
容错性 – Kafka 在多个节点间复制数据,确保即使某个 broker 故障,数据仍然可用。
持久性 – 消息持久化到磁盘,允许消费者在需要时重放事件。
支持事件驱动架构 – 它在微服务之间实现异步通信,非常适合现代云应用。
何时选择 Kafka:
高吞吐量、实时数据处理 – 非常适合日志处理、金融交易和 IoT 数据流。
微服务解耦 – Kafka 充当中介,允许微服务异步通信,无需直接依赖。
事件驱动系统 – 如果你的架构围绕对变化的反应(例如用户事件触发多个下游操作),Kafka 是个不错的选择。
可靠的消息传递和持久化 – 与可能丢弃消息的传统消息队列不同,Kafka 保留消息指定时间,确保持久性和可重放性。
可扩展性和容错性 – Kafka 的分布式特性允许它水平扩展,同时通过复制维护容错能力。
消息(Message) 是 Kafka 中数据的最小单位。
主题(Topic) 是生产者发送消息、消费者读取消息的逻辑通道。主题帮助对消息进行分类(例如日志、交易、订单)。
生产者(Producer) 是向主题发布消息的 Kafka 客户端。消息可以用三种方式发送:
Kafka 允许配置确认(ACKs)以平衡一致性和性能:
消息压缩与批处理 – Kafka 生产者在发送消息到 broker 之前可以对消息进行批处理和压缩。这提高吞吐量、减少磁盘使用但增加 CPU 开销。
Avro 序列化/反序列化 – 使用 Avro 而不是 JSON 需要预先定义模式,但提高性能并减少存储消耗。
Partition Kafka 主题被分成多个分区,允许并行处理和可扩展性。
Kafka 消费者采用轮询模型,意味着它们不断从 broker 请求数据,而不是 broker 主动推送数据给它们。
Partition 分配策略:
批大小配置 – 消费者可以定义每个轮询周期应检索多少条记录或多少数据。
消费者组 是一组消费者共同处理来自主题的消息。
Kafka 确保单个 partition 在一个组内仅由一个消费者消费,维持消息顺序。
当消费者读取消息时,它会更新其 offset(最后处理消息的位置)。
Broker 是存储消息、分配 offset 并处理客户端请求的 Kafka 服务器。
多个 broker 构成 Kafka 集群以实现可扩展性和容错。
Zookeeper 管理元数据、跟踪 broker 并处理 leader 选举。
不过,较新的 Kafka 版本正在努力消除 Zookeeper 依赖。
为了更好地理解 Kafka,让我们看一个简单例子,其中生产者向主题发送消息,两个不同的消费者分别处理这些消息:一个模拟电子邮件通知服务,另一个将消息存储到数据库。
services:
zookeeper:
image: confluentinc/cp-zookeeper:latest
container_name: zookeeper
restart: always
environment:
ZOOKEEPER_CLIENT_PORT: 2181
kafka:
image: confluentinc/cp-kafka:latest
container_name: kafka
restart: always
depends_on:
- zookeeper
ports:
- "9092:9092"
- "29092:29092"
environment:
KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092,PLAINTEXT_INTERNAL://kafka:29092
KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,PLAINTEXT_INTERNAL://0.0.0.0:29092
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_INTERNAL:PLAINTEXT
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
const { Kafka } = require("kafkajs");
const kafka = new Kafka({
clientId: "family-producer",
brokers: ["localhost:9092"],
});
const producer = kafka.producer();
async function sendMessage() {
await producer.connect();
console.log("🟢 Producer connected");
const message = {
id: Date.now(),
content: `Hi Mom! Time is ${new Date().getMinutes()}:${new Date().getSeconds()}`,
};
await producer.send({
topic: "family-topic",
messages: [{ value: JSON.stringify(message) }],
});
console.log(`📨 Sent: ${JSON.stringify(message)}`);
await producer.disconnect();
}
sendMessage();
const { Kafka } = require("kafkajs");
const kafka = new Kafka({
clientId: "family-email-consumer",
brokers: ["localhost:9092"],
});
const consumer = kafka.consumer({ groupId: "email-group" });
async function consumeMessages() {
await consumer.connect();
await consumer.subscribe({ topic: "family-topic", fromBeginning: true });
console.log("🟢 Email Consumer Connected");
await consumer.run({
eachMessage: async ({ message }) => {
const msg = JSON.parse(message.value.toString());
console.log(`📩 Notification Sent: "${msg.content}"`);
console.log(`📧 Email Sent: "${msg.content}" \n`);
},
});
}
consumeMessages();
const { Kafka } = require("kafkajs");
const kafka = new Kafka({
clientId: "family-db-consumer",
brokers: ["localhost:9092"],
});
const consumer = kafka.consumer({ groupId: "db-group" });
async function consumeMessages() {
await consumer.connect();
await consumer.subscribe({ topic: "family-topic", fromBeginning: true });
console.log("🟢 DB Consumer Connected");
await consumer.run({
eachMessage: async ({ message }) => {
const msg = JSON.parse(message.value.toString());
console.log(`💾 Storing message in DB: "${msg.content}" \n`);
},
});
}
consumeMessages();
Kafka 是一个强大的工具,已经改变了实时数据处理的方式。
然而,虽然它提供了令人难以置信的可扩展性和持久性,但评估它是否适合你的架构至关重要。
敬请期待!我将撰写一篇后续文章比较 Kafka 对比 Redis,以探索它们的用例和何时选择其中一个。🚀
Athreya aka Maneshwar
来源:部分图片来自这里:1
AI agent 写代码很快。但它们也会无声地删除逻辑、改变行为并引入 bug -- 都不告诉你。你通常在生产环境才发现。
git-lrc 解决了这个问题。它挂钩到 git 提交,在每个 diff 落地前进行审查。60 秒安装。完全免费。
欢迎任何反馈或贡献!它在线可用,代码开源,可供任何人使用。
免费、微型 AI 代码审查在提交时运行
| 🇩🇰 Dansk | 🇪🇸 Español | 🇮🇷 Farsi | 🇫🇮 Suomi | 🇯🇵 日本語 | 🇳🇴 Norsk | 🇵🇹 Português | 🇷🇺 Русский | 🇦🇱 Shqip | 🇨🇳 中文 | 🇮🇳 हिन्दी |
免费、微型 AI 代码审查在提交时运行
AI agent 写代码很快。但它们也会无声地删除逻辑、改变行为并引入 bug -- 都不告诉你。你通常在生产环境才发现。
git-lrc 解决了这个问题。它挂钩到 git 提交,在每个 diff 落地前进行审查。60 秒安装。完全免费。
查看 git-lrc 如何捕获严重的安全问题,如泄露的凭据、昂贵的云操作和日志语句中的敏感信息。
🤖 AI agent 无声地破坏事物。代码被删除。逻辑被改变。边界情况消失。你不会注意到直到生产。
🔍 在它发布前捕获。AI 驱动的内联注释准确显示改变了什么以及什么看起来有问题。
某些评论可能仅对登录访客可见。登录以查看所有评论。
如需进一步行动,你可以考虑屏蔽此人和/或举报滥用行为。