编辑
2026-06-04
java炒饭
00

Kafka 是一个 分布式、可持久化的消息队列系统,但它比传统消息队列更强:能存储大量数据、支持高吞吐、还能重复消费历史消息。


一、Kafka 的整体架构(先有个图景)

Kafka 集群由多个 Broker 组成(Broker 就是一台 Kafka 服务器)。
消息通过 Producer 发送到 Broker 的某个 Topic 里;
Consumer 从 Topic 拉取消息进行消费。

Producer1 ──┐ Producer2 ──┼──► Broker1 (Leader) ──┬──► Consumer Group A (成员1) Producer3 ──┘ Broker2 (Follower) └──► Consumer Group A (成员2) Broker3
  • ZooKeeper / KRaft:早期 Kafka 用 ZooKeeper 管理集群元数据(如 Broker 列表、Topic 配置);新版本用内置的 KRaft 替代 ZooKeeper,原理类似。

二、核心名词概念详解

1. Topic(主题)

  • 是什么:消息的分类容器。比如“订单消息”发到 order-topic,“用户行为”发到 user-action-topic
  • 特点:一个 Topic 可以有多个 Partition(分区)。分区是物理上的存储单元,每个分区是一个有序的消息日志文件。
  • 为什么分区:为了横向扩展。一个 Topic 可以分散存储在多个 Broker 上,通过多分区实现高吞吐(并行读写)。

2. Partition(分区)

  • 消息顺序:每个分区内部消息是严格有序的(FIFO),但跨分区不保证顺序。
  • 偏移量(Offset):每条消息在分区内的唯一序号,从 0 开始递增。消费者通过 Offset 来定位消息。
  • 副本(Replica):每个分区可以有多个副本(Leader + Follower),Leader 处理读写,Follower 同步数据,保证高可用。

3. Producer(生产者)

  • 负责发送消息到 Topic 的指定分区(可指定 key,相同 key 的消息去同一分区)。
  • 可以设置 acks 参数控制可靠性:
    • acks=0:不等待确认,最快,可能丢数据。
    • acks=1:Leader 写入成功即确认。
    • acks=-1/all:Leader 和所有 ISR(同步副本)都写入才确认,最可靠。

4. Consumer(消费者)

  • 从 Topic 的分区拉取消息。
  • 消费者需要指定 Consumer Group(消费者组)。

5. Consumer Group(消费者组)—— 这是 Kafka 最巧妙的设计

  • 定义:一组具有相同 group.id 的 Consumer 实例。
  • 核心规则:一个分区内的消息,只能被同一个消费者组内的 一个 消费者消费。
  • 作用
    • 负载均衡:如果组内有多个消费者,Kafka 会自动将分区分配给它们,实现并行消费。
    • 容错:如果某个消费者挂了,它负责的分区会重新分配给组内其他消费者。
    • 消息不重复消费:不同消费者组对同一条消息是“独立”的——两个不同组可以消费同一条消息,就像两个独立的订阅者。

举例:Topic order-topic 有 3 个分区(P0, P1, P2)。

  • 消费者组 G1 有 3 个消费者:每人分一个分区,并行消费,每条消息被消费 1 次
  • 消费者组 G2 有 2 个消费者:其中一个消费者分到 2 个分区,另一个分到 1 个分区,消息在 G2 内也只消费 1 次。
  • 如果 G1 和 G2 同时消费,同一条消息会被两个组各消费一次 —— 不同组的消费是隔离的

6. Broker(代理节点)

  • 一台 Kafka 服务器就是一个 Broker,负责存储消息、处理读写请求。
  • 多个 Broker 组成集群。一个 Broker 可以承载多个分区的 Leader 或 Follower。

7. ISR(In-Sync Replicas)

  • 与 Leader 保持同步的副本集合。如果 Follower 复制落后太多,会被踢出 ISR。
  • 只有 ISR 中的副本才有资格被选为新 Leader。

8. Offset(偏移量)

  • 每个分区内消息的序号(long 型)。消费者需要提交自己已经消费到的 Offset,以便故障恢复后继续消费。
  • Kafka 将每个消费者组的 Offset 保存在一个内部 Topic __consumer_offsets 中。

三、Topic 与消费者组的关系(重点)

特性TopicConsumer Group
角色消息的存放地消费者的组织单位
关键属性分区数、副本数group.id
消息流向Producer → Topic → Consumer组内消费者共同承担分区消费
消息消费范围每个分区只能被同一个组内的一个消费者消费不同组可以独立消费全部消息

一个经典的提问:如果消费者组内的消费者数量 > Topic 的分区数量,会怎样?
→ 多余的消费者将空闲(没有分区可分配)。因此一般设置消费者数 ≤ 分区数。


四、消息流转完整过程(结合概念)

  1. 创建 Topic:指定分区数(例如 3)、副本数(例如 2)。
  2. Producer 发送消息:消息被追加到某个分区的 Leader 副本末尾。
  3. Broker 同步:Follower 从 Leader 拉取消息,保持同步。
  4. Consumer 启动:加入指定的 Consumer Group,Kafka 根据分区分配策略(如 Range、RoundRobin)分配分区给消费者。
  5. Consumer 拉取消息:不断拉取自己分配到的分区,处理消息,并定期提交 Offset(自动或手动)。
  6. 重平衡(Rebalance):当消费者组内成员变动(新增、退出、崩溃),Kafka 会重新分配分区,这个过程叫 Rebalance。期间该组会短暂停止消费。

五、为什么 Kafka 这么强?(基于以上架构的优势)

  • 高吞吐:分区并行 + 批量读写 + 顺序 IO。
  • 持久化:消息直接刷盘,且可配置保留时间(默认 7 天)。
  • 高可用:副本机制 + ISR + 自动 Leader 选举。
  • 可回溯消费:因为消息不删除(按时间/容量删除),消费者可以重置 Offset 到过去某个时刻重新消费。
  • 解耦:不同消费者组独立消费,互不影响。

六、一张简表总结

概念一句话解释
Topic逻辑上的消息分类
PartitionTopic 的物理分片,提供顺序和并行能力
Producer发消息的人
Consumer收消息的人
Consumer Group一组消费者,共同消费一个 Topic,彼此分担负载
BrokerKafka 服务器节点
Offset分区内消息的序号
ISR与 Leader 保持同步的副本集合

编辑
2026-06-04
java炒饭
00

下面解析雪花算法(Snowflake) 的核心原理,以及它的优势与局限。

这是Twitter在2010年开源的分布式ID生成算法,旨在解决分布式系统中需要全局唯一、趋势递增、高性能的ID生成问题。它的输出是一个64位的long型整数,结构清晰,非常经典。


一、原理:64位ID的精细划分

雪花算法的64位(bit)被划分为以下几个部分(从左到右,高位到低位):

位数字段名含义与作用
1位符号位固定为0。因为ID是正数,最高位为0。
41位时间戳记录当前时间与某个起始时间的差值(毫秒级)。可支持约69年。
10位工作机器ID用于区分不同的节点/机器,最多支持1024个节点(2^10)。实际常拆分为5位数据中心ID + 5位机器ID。
12位序列号同一毫秒内生成的不同ID的序号,支持每毫秒每节点最多4096个ID(2^12)。

整体结构示意图:

0 | 41位时间戳(毫秒) | 10位机器ID | 12位序列号

生成流程(简化):

  1. 获取当前时间戳(毫秒)。
  2. 如果当前时间戳 等于 上一次生成ID的时间戳:
    • 序列号自增1。
    • 若序列号超过4095,则等待下一毫秒,并将序列号重置为0。
  3. 如果当前时间戳 大于 上一次的时间戳:
    • 序列号重置为0。
  4. 如果当前时间戳 小于 上一次的时间戳(即时钟回拨),则算法通常停止生成或等待。
  5. 将各部分按位或(左移后相加)组合成一个64位long整数。

示例(假设起始时间为2010-11-04 09:42:54.657):

  • 时间戳差值:某个毫秒值
  • 机器ID:1
  • 序列号:5 → 最终得到一个类似 1456192189189701632 的ID。

由于时间戳在高位,生成的ID趋势上是递增的(但不严格连续,同毫秒内序列号连续)。


二、优势

  1. 全局唯一,分布式友好
    只需确保每个节点拥有唯一的机器ID,就可以独立生成ID,无需中心化协调。非常适合微服务、分库分表场景。

  2. 趋势递增,利于数据库索引
    生成的ID在毫秒级时间上递增(同毫秒内序列号递增)。对数据库B+树索引非常友好,减少页分裂,写入性能高。

  3. 高性能
    生成一个ID仅需几次位运算和内存操作,无网络IO、无锁竞争(单节点内常用atomic实现)。单机轻松达到数十万甚至百万级QPS

  4. 灵活性
    位数分配可按需调整。例如增加机器位或序列号位,适应不同规模集群。

  5. 无外部依赖
    算法只依赖当前机器的时间戳,不依赖数据库、ZooKeeper、Redis等第三方服务,轻量可靠。


三、劣势与挑战

  1. 强依赖机器时钟
    如果节点的系统时钟发生回拨(向前调整时间),那么有可能生成重复ID或导致新ID小于已生成的ID。

    • 常见原因:NTP同步、人为修改、网络延迟等。
    • 处理方案
      • 检测到回拨时:直接抛出异常阻塞等待使用备用时钟源
      • 更健壮的方案(如百度UidGenerator):使用未来时间生成策略 + 回拨容忍。
  2. 机器ID的分配与管理问题
    每个节点需要独立的机器ID(10位)。如何安全、动态地分配这1024个ID而不冲突?

    • 简单方案:手动配置。
    • 常见方案:使用ZooKeeper、etcd或数据库自增表做中心化分配。
    • 在容器或弹性伸缩场景下,自动分配ID是个难题。
  3. 时间跨度有限
    41位时间戳(毫秒级,起始点固定)可用约69年。如果服务持续运行超过69年,需要变更起始时间或扩展位数(通常业务系统远达不到这个时限)。

  4. 序列号容量限制
    每毫秒每节点最多4096个ID。若瞬时并发超过这个量,可能产生少量等待(自旋到下一毫秒)。但在绝大多数场景下足够。

  5. 无法保证全局严格递增
    “趋势递增”≠“单调递增”。不同节点由于时钟差异,可能会出现后生成的ID比先生成的ID数值小的情况(跨节点比较)。如果要求所有ID全局严格递增,雪花算法不满足。


编辑
2026-01-10
java炒饭
00

该文章已加密,点击 阅读全文 并输入密码后方可查看。

编辑
2026-01-09
redis
00

单节点的redis并发能力有限,因此我们需要搭建主从集群,实现读写分离。一般是主节点进行写,从节点进行读。

主从同步:

  • 全量同步:从结点第一次与主节点同步数据时使用的方案。从节点向主结点发送同步请求,携带参数replicationId,offset,主节点会根据replicationId来判断是否是第一次同步,如果是第一次同步(Id不一致),则主节点会把自己的replicationId和offset发给从节点,同时主节点执行bgsave生成rdb文件。在生成的同时,会开启一个缓冲区记录该期间收到的所有写命令,最后把rdb和这个缓冲区的信息(类似日志)一起发送给从节点。
  • 增量同步:还是从节点发送请求,主节点判断是不是第一次请求,不是则获取从节点的offset,并把从结点的offset和自身offset之间的数据同步给从节点。

高可用:

哨兵模式(sentinel):

  • 本质自己也是一个redis服务
  • 每1s向集群的每个实例发送ping请求。如果在规定时间收不到响应,认为该结点主观下线。
  • 当指定数量的哨兵认为某个结点都下线了,则认为该节点客观下线。这个数量一般不低于哨兵数量的一半。
  • 当主节点下线了,哨兵会推选一个新的主节点。选主规则:
    • 如果该结点与原来的主节点断开连接时间超过指定值,则不纳入考虑。
    • 根据结点的优先级来判断,优先选择优先级高的,可在配置文件中设置优先级。
    • 优先级相同时选择offset大的节点
    • 以上条件都相同时选择运行ID最小的节点

脑裂:

  • 由于网络原因,哨兵监测不到主结点,认为主节点挂了,实际上主节点并没挂,此时客户端还在与原来的主节点发送请求,然而哨兵却选择了一个新的主节点,原来的主节点变成了从节点,在同步数据的时候,期间所写的数据就失效了。
  • 解决方案:可以通过配置min-slaves-to-write(最少从节点数)和min-slaves-max-lag(最大延迟秒数)。当主节点的从节点数量少于配置值,或者从节点的延迟时间超过配置值时,主节点会拒绝写入请求,从而避免数据不一致。这是一种 “宁可拒绝写入,也要保证数据一致性” 的取舍。在脑裂期间,被孤立的旧主节点会提供高可用性。
编辑
2026-01-09
redis
00

redis实现的分布式锁:

  • 使用setnx命令时,需要设置 ttl,防止系统故障导致锁无法释放。
  • 自己实现的分布式锁的缺陷:我们并不知道准确的业务执行时间,因此这个过期时间不好控制。
  • 不可重入

因此我们使用第三方工具redisson:

  • 提供看门狗(WatchDog),一个线程获取锁成功之后,WatchDog会给持有锁的线程续期(默认每隔10s续期)
  • 可重入,底层采用了一个hash结构,用线程id和该锁锁的次数作为依据,如果发现锁已经被获取了,但是是当前线程获取的,我们就可以再次获得锁,并把次数加1。如果发现这个线程id不是自己的,则无法获取锁。释放锁的时候让次数减1即可。
  • 不能做到主从的强一致性,如果需要,可以使用zookeeper实现的分布式锁。
  • 底层还是setnx和lua