安全体系:认证、授权与加密
上一篇分析了性能模型。这一篇进入生产环境的另一个关键维度:安全。 Kafka 早期版本没有内置安全机制——客户端连上 broker 就能读写任意 topic。自 0.9 版本起,Kafka 逐步引入了认证、授权和加密三层安全能力。这三层各自解决一个问题:认证确认身份,授权控制权限,加密保护传输通道。 本文只抓一个问题:Kafka 的安全体系如何分层工作,每一层的配置和机制是什么。 三层安全模型 1234567891011Client Broker | | |--- TLS handshake ------------>| Layer 3: Encryption |<-- certificate exchange ------| (通道加密) | | |--- SASL authentication ----->| Layer 1: Authentication |<...
性能模型:吞吐、延迟与调优思路
前面的文章覆盖了 Kafka 的核心机制和生态组件。这一篇回到工程视角:Kafka 的性能模型到底长什么样。 性能调优容易变成一张参数清单。更有效的方式是先建立一个吞吐-延迟的基本模型,再用这个模型解释每个参数的作用方向和代价。 本文只抓一个问题:Kafka 的吞吐和延迟分别由哪些因素决定,调节一个参数时另一端会发生什么变化。 性能基础:四个底层机制 Kafka 的高吞吐不是来自某个单一优化,而是四个机制叠加的结果。 1234567891011Producer Broker Consumer | | | |-- batch N records ---->| | | (linger.ms window) |-- sequential append --> | | | (page cach...
Schema Registry 与数据治理:给消息加上契约
上一篇解决了数据管道的标准化。这一篇解决管道里的数据该长什么样——消息的 schema 治理。 消息格式的约定容易被忽视。producer 和 consumer 对消息格式的理解仅靠代码约定和口头沟通时,任何一方的字段变更都可能在运行时导致反序列化失败。Schema Registry 的目标是把这种隐式约定变成显式的、可验证的契约。 本文只抓一个问题:Schema Registry 如何通过版本化 schema 和兼容性检查,在不停机的情况下实现消息格式的安全演进。 没有 Schema 治理时的问题 12345678910Producer v1 Consumer v1{name: "alice"} → 解析 {name} ✓Producer v2 (加字段) Consumer v1 (未升级){name: "alice", → 解析 {name, age} ? ...
Kafka Connect:标准化的数据管道
上一篇看了 Kafka 内置的流处理库。这一篇转向数据集成:怎么把外部系统的数据可靠地搬进搬出 Kafka。 数据集成容易退化成"给每个数据源写一个 producer,给每个数据目标写一个 consumer"。当数据源和目标各有几十个时,组合爆炸。Kafka Connect 的目标是把这个问题标准化:用统一的框架和插件机制,把数据管道的公共逻辑(offset 追踪、序列化、容错、分布式任务分配)收到框架层,让开发者只关注"怎么从特定系统读数据"和"怎么往特定系统写数据"。 本文只抓一个问题:Kafka Connect 如何用标准化的 connector/task 模型实现可靠的数据管道。 架构模型 123456789101112 Kafka Connect Cluster┌────────────────────────────┐│ Worker 1 Worker 2 ││ ┌────────┐ ┌────────┐ ││ │Task A-0│ │Task A-1│ │ ← So...
Kafka Streams:在日志之上构建流处理
前十篇覆盖了 Kafka 作为消息基础设施的核心机制。这一篇开始看 Kafka 的生态扩展——第一个是内置的流处理库 Kafka Streams。 流处理容易被误解成"又一套需要独立部署的集群"。Kafka Streams 的核心定位不是集群,而是一个嵌入应用 JVM 的 Java 库。它把 Kafka topic 当作输入和输出,在应用进程内完成流式计算。 本文只抓一个问题:Kafka Streams 怎样在追加日志之上构建出有状态的流处理。 架构定位:库,不是集群 12345678910111213141516┌─────────────────────────────────────────┐│ Application JVM ││ ┌───────────────────────────────────┐ ││ │ Kafka Streams Library │ ││ │ ┌─────────┐ ┌───────────────┐ │ ││ │ │ Topol...
日志压缩:把 topic 当 KV 表用
前面的文章中,日志的保留策略一直是基于时间或大小的删除。这一篇介绍另一种保留策略:日志压缩。 日志压缩容易被理解成"压缩存储空间"。更准确的说法是:日志压缩是一种按 key 去重的保留策略,对每个 key 只保留最新的 value,把一个 append-only 的日志变成一个可以反映最新状态的 KV 快照。 本文只抓一个问题:日志压缩的语义、执行机制和适用场景。 两种保留策略 Kafka 的日志保留由 topic 级别的 cleanup.policy 参数控制,有三种取值: 123cleanup.policy=delete 基于时间或大小删除整个 segmentcleanup.policy=compact 按 key 去重,保留每个 key 的最新 valuecleanup.policy=compact,delete 先压缩去重,再按时间/大小删除 delete 策略是前面几篇文章一直在讨论的默认行为:segment 文件超过 retention.ms 或 retention.bytes 后被整体删除,不关心消息的 key 是什么。 comp...
Exactly-Once 与事务:跨 partition 的原子写入
前面在第 04 篇解决了单 partition 内的消息去重。这一篇解决更难的问题:一批消息写到多个 partition,要么全部可见,要么全部不可见。 幂等 Producer 的序列号去重是 per-PID、per-partition 的。它不解决两个场景:一是 Producer 重启后 PID 改变导致无法关联之前的写入;二是一个业务操作涉及多个 partition 的写入,部分成功部分失败时没有回滚机制。Kafka 的事务机制(Transactional API)正是为了解决这两个问题而设计的。 本文只抓一个问题:Kafka 事务是怎么通过 transactional.id、Transaction Coordinator 和两阶段提交协议实现跨 partition 原子写入的。 为什么幂等不够 用一个具体场景说明。一个 stream processing 应用从 input topic 消费消息,处理后写入 output topic,同时提交 consumer offset。这三个动作涉及不同的 partition: 12input-topic (partition-0) ...
Controller 与 KRaft:从 ZooKeeper 到内置共识
上一篇看到了副本机制的细节。这一篇解决元数据层的核心问题:谁来决定哪个 replica 是 leader,以及这个决定怎么达成共识。 这个问题容易和副本选举混淆。副本选举本身只是一个结果——“partition-0 的 leader 从 broker-1 变成 broker-2”。真正的问题是:谁有权做出这个决定,这个决定如何传播到集群中的每个 broker,以及当做出决定的那个节点自己挂了会发生什么。 本文只抓一个问题:Kafka 的元数据管理从外部依赖(ZooKeeper)演化到内置共识(KRaft)的完整路径。 ZooKeeper 模式下的 Controller 在 ZooKeeper 模式下,Kafka 集群中有且仅有一个 broker 担任 controller 角色。Controller 是通过 ZooKeeper 的临时节点(ephemeral node)竞选产生的: 123456789101112 ZooKeeper Ensemble ┌───────────────────┐ │ /co...
副本与 ISR:高可用的代价和折中
前六篇走完了生产和消费两端。这一篇进入 Kafka 的高可用核心:副本机制。 副本机制容易被简化成"多写几份"。更准确的说法是:Kafka 在每个 partition 级别维护一组副本,用 ISR(In-Sync Replicas)动态集合代替固定多数派投票,在可用性和持久性之间做出了一系列显式的折中。 本文只抓一个问题:一条消息从被 leader 接收到被所有 ISR 副本确认,中间经历了哪些步骤,以及这些步骤中的每一个参数如何影响数据安全。 副本模型 Kafka 的副本模型是 per-partition 的 leader-follower 结构: 12345678910111213141516171819202122Producer │ ▼┌──────────────┐│ Partition-0 ││ Leader (B-0) │ ◄── 所有读写都经过 leader│ [0][1][2][3] │└──────┬───────┘ │ FetchRequest (replicaId=1) ▼┌─────────...
Offset 管理:提交、重置与消费语义
上一篇解决了 partition 怎么分给 consumer。这一篇解决分到之后的核心问题:consumer 读到哪里了,怎么记住。 这个"记住"机制看起来简单——不过是记录一个数字(offset)。但"什么时候记"和"怎么记"的选择,直接决定了消息会不会丢、会不会重复处理。auto-commit 的默认行为经常被误解为"自动保证 exactly-once",实际上它给出的是 at-least-once 甚至更弱的语义。 本文只抓一个问题:offset 的提交策略如何影响消费语义。 三个 offset 理解 offset 管理需要区分三个位置: 123456789Partition Log:┌────┬────┬────┬────┬────┬────┬────┬────┬─────┐│ 0 │ 1 │ 2 │ 3 │ 4 │ 5 │ 6 │ 7 │ ... │└────┴────┴────┴────┴────┴────┴────┴────┴─────┘ ...
