Kafka 消息队列实战:个人站的异步解耦,Docker 部署、消费积压与数据不丢的五个前提

为什么个人站需要 Kafka,什么时候不需要

先泼一盆冷水:大多数个人站长根本用不上 Kafka。Kafka 是分布式消息队列,它的设计目标是「每秒几十万条消息、多个消费者、可回溯的日志流」。如果你只是想给注册成功发一封欢迎邮件,或者写个定时任务,用 Redis 的 list、甚至直接一个 cron + 数据库轮询就够了,上 Kafka 属于拿高射炮打蚊子。

但有几类场景,一旦你的站开始碰上,就会立刻感受到「没有消息队列真难受」:

  • 用户注册后要发邮件、写积分、加统计,如果全放在注册接口里同步做,一个邮件服务器卡顿,注册接口就跟着超时;
  • 网站日志、埋点数据需要异步落盘和分析,直接写数据库会把数据库连接占满;
  • 多个下游系统(统计、告警、报表)都需要同一份原始数据,同步「扇出」写多份既慢又容易漏;
  • 偶尔有突发流量,希望先把请求堆进队列,后端慢慢消费,而不是直接被压垮。

这些问题的本质都是「生产」和「消费」的速度不匹配,需要一个缓冲区把两者解耦。Kafka 相比 Redis 这类简易队列的优势在于:消息可以持久化到磁盘、可以重复消费、支持多个独立的消费组各自读取全量数据。对个人站来说,前两点尤其重要——服务器重启不丢消息,出问题时还能把历史消息倒出来重放。

Kafka 的核心概念速览

Kafka 术语不少,但真正需要理解的只有四个:

Topic(主题)是消息的类别,相当于一个队列的名字。生产者往 topic 里发消息,消费者从 topic 里读消息。

Partition(分区)是 topic 的物理切分。一个 topic 可以分成多个分区,分布在不同的 broker 上。分区是 Kafka 并行度的基本单位:一个分区在任一瞬间只能被同一个消费组里的一个消费者读取。因此,分区数直接决定了你的最大并发消费能力。

Offset(偏移量)是消息在分区内的唯一编号,单调递增。消费者通过记录自己读到了哪个 offset 来标识「消费到哪了」。这个特性让 Kafka 支持重复消费——把 offset 回退,就能重放历史消息。

Consumer Group(消费组)是一组协同消费的消费者。同一个组内,每个分区只会分给一个消费者,从而实现负载均衡;不同组之间互不影响,各自都能读到全量消息。

还有一个关键点:Kafka 的消息默认保留一段时间(比如 7 天)或一定体积,到期才删除,而不是「消费完就删」。这使得它更像一个可回溯的日志系统,而不是传统意义上「读完就没了」的队列。

用 Docker Compose 单机部署 Kafka

个人站没必要上集群,单机一个 broker 完全够用。这里用 Docker Compose 部署一个 KRaft 模式(不需要 ZooKeeper)的 Kafka:

services:
  kafka:
    image: bitnami/kafka:3.7
    container_name: kafka
    restart: unless-stopped
    ports:
      # 只监听内网,绝不暴露公网
      - "10.0.0.11:9092:9092"
    environment:
      KAFKA_CFG_NODE_ID: "1"
      KAFKA_CFG_PROCESS_ROLES: "broker,controller"
      KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: "1@kafka:9093"
      KAFKA_CFG_LISTENERS: "PLAINTEXT://:9092,CONTROLLER://:9093"
      KAFKA_CFG_ADVERTISED_LISTENERS: "PLAINTEXT://10.0.0.11:9092"
      KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: "CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT"
      KAFKA_CFG_CONTROLLER_LISTENER_NAMES: "CONTROLLER"
      KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE: "false"
      KAFKA_CFG_LOG_RETENTION_HOURS: "168"
      KAFKA_HEAP_OPTS: "-Xmx512m -Xms256m"

几个关键点值得说明:

ADVERTISED_LISTENERS 告诉客户端「连我应该用哪个地址」。这个值必须是客户端能访问到的地址,配错了会连不上,是新手最常见的坑。

AUTO_CREATE_TOPICS_ENABLE: "false" 建议关掉自动建 topic,改成手动建,避免拼错主题名时悄悄创建一个新 topic 却没人消费。

KAFKA_HEAP_OPTS 限制 JVM 堆内存。Kafka 是 Java 应用,默认可能占几个 G,低配 VPS 一定要限制,否则会被 OOM Killer 干掉。

启动后进入容器创建 topic:

docker exec -it kafka bash
kafka-topics.sh --bootstrap-server localhost:9092 \
  --create --topic site-events \
  --partitions 3 --replication-factor 1
kafka-topics.sh --bootstrap-server localhost:9092 --list

单机只有一个 broker,所以副本数只能是 1。分区数可以给 3,为将来横向扩展消费端留余地。

用什么语言写生产者和消费者

Kafka 的客户端几乎覆盖所有语言。个人站最常用的组合是:Web 端(PHP)做生产者,独立的常驻脚本(Python)做消费者。先看 PHP 生产者(需要 rdkafka 扩展):

<?php
$conf = new RdKafka\Producer();
$conf->addBrokers("10.0.0.11:9092");

$topic = $conf->newTopic("site-events");
$payload = json_encode([
    "type" => "user_register",
    "uid"  => 1024,
    "ts"   => time(),
], JSON_UNESCAPED_UNICODE);

$topic->produce(RD_KAFKA_PARTITION_UA, 0, $payload);
$conf->poll(0);   // 触发实际发送
$conf->flush(1000); // 确保刷出,最多等 1 秒
?>

生产者的关键原则是异步、不阻塞主流程。注册接口把消息丢进 Kafka 就返回,至于发邮件、算积分交给消费者慢慢做。如果 Kafka 短暂不可用,要能降级(比如落本地文件补发),不能让注册直接失败。

再看 Python 消费者:

from kafka import KafkaConsumer
import json

consumer = KafkaConsumer(
    "site-events",
    bootstrap_servers=["10.0.0.11:9092"],
    group_id="site-worker",
    auto_offset_reset="earliest",
    enable_auto_commit=False,   # 手动提交,处理成功才提交
    value_deserializer=lambda m: json.loads(m.decode("utf-8")),
)

for msg in consumer:
    try:
        handle(msg.value)          # 业务处理
        consumer.commit()          # 处理成功才提交 offset
    except Exception as e:
        print("处理失败,不提交,等待重试:", e)

enable_auto_commit=False 是非常关键的一行。自动提交会在「消息收到」时就提交 offset,如果处理过程中脚本崩溃,这条消息就永远丢了。手动提交能保证至少消费一次(at-least-once),代价是可能重复处理,所以消费逻辑要写成幂等的。

消费积压怎么看、怎么处理

消息队列最怕的不是慢,而是积压——生产者发得比消费者处理得快,lag(滞后量)越来越大。查看每个消费组的积压情况:

kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --describe --group site-worker

输出里的 LAG 列就是积压量。判断原则很简单:LAG 稳定在一个小数字说明消费跟得上;LAG 持续增长说明消费者能力不足。解决办法按优先级:

  • 提高消费并行度:增加消费者实例,但记住并行度上限是分区数——3 个分区最多 3 个有效消费者,第 4 个只会空转。所以建 topic 时要预估好分区数。
  • 加快单条处理:把消费者里耗时的外部调用(发邮件、调 API)改成批量处理,或者放进另一个队列。
  • 临时放宽:如果积压是一次性的(比如某次活动),可以临时多开几个消费者把积压清掉,事后恢复。

反过来,如果 lag 长期为 0 甚至消费者空闲,说明资源浪费,可以适当降低消费者数量。

数据不丢的几个前提

Kafka 号称持久化,但默认配置下仍然可能丢消息,要保证不丢需要理解这几件事:

生产端。把 acks 设为 all,意味着消息要等所有同步副本确认才认为发送成功。单机 broker 场景下,这一项加上适当的 retries,能避免「发出去了但 broker 其实没存下」。

消费端。坚持手动提交 offset,处理成功再提交。宁可重复(幂等),不要丢失。

保留策略。默认保留 7 天,如果消费彻底停了超过 7 天,数据会被删掉。重要数据要么拉长 log.retention.hours,要么让消费者及时跟进。

磁盘。Kafka 的数据全在磁盘上,务必单独挂盘,并监控磁盘使用率。磁盘写满会让 broker 直接挂掉。

五个必踩的坑

坑一:把 9092 暴露到公网。Kafka 默认没有任何认证,暴露到公网等于把整个消息系统敞开给别人。永远只监听内网地址,需要跨机访问就走 WireGuard 之类的内网隧道。要用 SASL 认证就把配置写全,别只改一半。

坑二:分区数建少了。分区数是「不可逆」的——虽然 Kafka 支持增加分区,但增加后同一个 key 的分区归属会变,破坏 key 的有序性。所以建 topic 时按最大预期并发的两三倍来给分区。

坑三:把 Kafka 当数据库用。Kafka 是流式管道,不是查询存储。要长期保存和分析的数据,还是应该消费出来后写进 ClickHouse、MySQL 或对象存储。

坑四:忘了消费者组会有 rebalance。当消费者上下线时,分区会在组内重新分配,期间会有短暂停顿。如果消费者本身处理很慢、又频繁重启,会反复触发 rebalance,整体吞吐崩掉。办法是设置合理的 session.timeout.ms 和心跳,并避免频繁重启。

坑五:不监控磁盘和 lag。Kafka 出问题往往是「悄悄积压」而不是「直接报错」。给自己的站加两个最简单的监控:消费组 lag 和 Kafka 数据目录的磁盘使用率,超阈值就发邮件告警。这两个指标能提前几小时甚至几天发现问题。

给个人站的落地建议

如果你决定引入 Kafka,建议按这个顺序落地:

  • 先只把「一个」明确的异步场景迁进去,比如注册后的所有副作用(邮件、积分、统计)统一走一个 topic,验证整条链路可靠;
  • topic 命名要带业务前缀和版本,比如 site-events-v1、site-emails-v1,方便将来平滑升级;
  • 所有消费逻辑写成幂等的,因为 at-least-once 一定会有重复;
  • 给消费者写 systemd 服务或 Supervisor 守护,崩溃自动拉起;
  • 上线后盯着 lag,如果长期为 0,说明这套东西对你可能是过度设计,可以简化掉。

Kafka 不是必需品,但它是一把很趁手的工具。真正决定成败的不是能不能装起来,而是你有没有想清楚「哪些操作该异步、消息丢了或重复了怎么办」。把这两个问题回答清楚,再动手部署,Kafka 才能成为你站点的稳定支柱,而不是又一个吃完内存又没人维护的玩具。

Last modification:October 4th, 2026 at 10:24 pm

Leave a Comment