支付事件、日志采集、订单状态和实时数据一旦由多个服务同时生产与消费,直接调用数据库或 HTTP 接口很容易形成耦合。Apache Kafka 把消息写入可分区、可回放的追加日志,让生产者和消费者按各自节奏工作,适合事件流、数据管道、日志汇聚和需要保留消费历史的业务。
本文以 Apache Kafka 4.3.1、KRaft、Docker Compose 和 Ubuntu 24.04 为基准,在 VPS 上部署一套单节点 Kafka。方案使用官方镜像、SASL/PLAIN + TLS、ACL、持久化目录、健康检查和固定版本,并覆盖容量规划、监控、冷备份、升级与常见故障。版本与官方文档按 2026 年 9 月 15 日核对;上线前仍应重新检查补丁版本与升级说明。
Kafka 的核心优势不是“把任务放到队列里”这么简单,而是对事件进行分区、持久化和重复消费:
- 订单、支付、库存等服务通过事件解耦;
- 日志或埋点先进入 Kafka,再由多个下游分别消费;
- 消费者可以保留 offset,故障后从指定位置继续;
- 同一个 Topic 可由多个 Consumer Group 独立读取;
- 通过分区并行提升吞吐,并按保留策略保存一段时间的数据。
如果只需要低延迟任务分发、复杂路由、逐条确认和死信队列,RabbitMQ 更直接。如果需要事件回放、流处理或一份数据供多个下游独立消费,Kafka 通常更合适。
| 架构 | 优点 | 局限 | 适用场景 |
|---|---|---|---|
| 单节点 KRaft | 成本低、部署和排障简单 | VPS、磁盘或进程故障都会中断服务;副本数只能为 1 | 开发、测试、可接受停机恢复的小型业务 |
| 同一 VPS 三容器 | 能验证集群配置 | 仍共享主机、磁盘和网络,不具备故障域隔离 | 仅用于实验 |
| 3 个独立节点 | Topic 可设 3 副本,允许 1 个节点故障后维持多数派 | VPS 成本与运维复杂度更高 | 关键生产业务 |
本文的 combined mode 让一个进程同时承担 broker 与 controller,适合单节点。它不是高可用架构。关键业务至少应使用 3 个独立故障域,并采用 replication.factor=3、min.insync.replicas=2 和生产者 acks=all;官方说明这组设置要求多数副本持久化后才确认写入。
Kafka 4.x 使用 KRaft,不再需要 ZooKeeper。新集群还支持动态 controller quorum,但单节点 Docker 示例使用静态 controller.quorum.voters 更容易部署与恢复。将来扩为关键集群时,应重新设计 controller 与 broker 角色,不要简单复制同一个数据目录。
Kafka 很依赖磁盘顺序写、页缓存和网络吞吐。容量评估至少要看峰值写入量、消息大小、分区数、保留时间、副本数和消费者回放速度。
| 使用规模 | VPS 建议 | 说明 |
|---|---|---|
| 开发验证 | 2 vCPU、4GB 内存、50GB SSD | Heap 1GB,数据可随时重建 |
| 小型单节点 | 4 vCPU、8GB 内存、100GB 以上 NVMe | Heap 2GB,给 Linux 页缓存保留空间 |
| 关键生产集群 | 3 台起,每台 4–8 vCPU、8–16GB 内存、独立 NVMe | 先按真实消息大小、吞吐和保留时间压测 |
不要把 JVM Heap 吃满整台 VPS。Kafka 会利用操作系统页缓存读取日志;8GB 主机给 Heap 2GB 通常比给 Heap 6GB 更稳。数据盘也不要与高 I/O 数据库共用。
本文只映射一个客户端端口:
| 端口 | 用途 | 暴露策略 |
|---|---|---|
9094/TCP | 客户端 SASL_SSL | 默认绑定 127.0.0.1;远程连接改绑 WireGuard 私网 IP |
19094/TCP | 容器内 broker 通信 | 不映射到宿主机 |
29093/TCP | KRaft controller | 不映射到宿主机 |
不要把 9092 或 controller 端口直接暴露公网。即使 9094 已启用认证与 TLS,也应在云安全组和主机防火墙中只允许业务服务器来源。可先按数据库与中间件公网暴露安全清单收紧边界。
以下命令以 Ubuntu 24.04 为例。服务器已有 Docker Compose v2 时,不要重复安装第三方脚本。
sudo apt update
sudo apt full-upgrade -y
sudo apt install -y ca-certificates curl openssl openjdk-21-jre-headless
docker --version
docker compose version
keytool -help >/dev/null
生产环境还应配置容器日志轮转、资源限制和回滚流程,可参考 Docker Compose 生产环境配置清单。
sudo install -d -m 0750 -o "$USER" -g "$USER" /opt/kafka
cd /opt/kafka
install -d -m 0750 data backups
install -d -m 0700 secrets pki
umask 077
docker run --rm --entrypoint /opt/kafka/bin/kafka-storage.sh \
apache/kafka:4.3.1 random-uuid
记下最后一条命令返回的 Cluster ID。创建 .env,把域名和私网监听地址换成自己的值:
CLUSTER_ID=请替换为上一步输出
KAFKA_HOST=mq.example.com
KAFKA_BIND_IP=127.0.0.1
默认只允许本机连接。如果应用通过 WireGuard 地址 10.66.0.1 访问,就把 KAFKA_BIND_IP 改为该地址,并让 mq.example.com 在业务服务器上解析到这个私网 IP。
生成十六进制密码,避免特殊字符破坏 JAAS 语法:
ADMIN_PASSWORD="$(openssl rand -hex 32)"
APP_PASSWORD="$(openssl rand -hex 32)"
STORE_PASSWORD="$(openssl rand -hex 24)"
printf '%s\n' "$STORE_PASSWORD" > secrets/kafka_keystore_creds
printf '%s\n' "$STORE_PASSWORD" > secrets/kafka_ssl_key_creds
printf '%s\n' "$STORE_PASSWORD" > secrets/kafka_truststore_creds
printf 'KafkaServer {\n org.apache.kafka.common.security.plain.PlainLoginModule required\n username="admin"\n password="%s"\n user_admin="%s"\n user_app="%s";\n};\n' \
"$ADMIN_PASSWORD" "$ADMIN_PASSWORD" "$APP_PASSWORD" \
> secrets/broker_jaas.conf
admin 只用于运维,业务程序使用 app。SASL/PLAIN 的密码不会在协议层自行加密,所以必须与 TLS 一起使用;不要把 JAAS、客户端配置或 .env 提交到 Git。
下面用内部 CA 为 kafka 容器名和 mq.example.com 签发服务器证书:
cd /opt/kafka/pki
openssl genrsa -out ca.key 4096
openssl req -x509 -new -sha256 -days 3650 \
-key ca.key \
-subj '/CN=Kafka Internal CA' \
-out ca.crt
openssl genrsa -out broker.key 4096
openssl req -new -sha256 \
-key broker.key \
-subj '/CN=mq.example.com' \
-out broker.csr
printf '%s\n' \
'subjectAltName=DNS:kafka,DNS:mq.example.com' \
'extendedKeyUsage=serverAuth' > broker.ext
openssl x509 -req -sha256 -days 825 \
-in broker.csr \
-CA ca.crt \
-CAkey ca.key \
-CAcreateserial \
-extfile broker.ext \
-out broker.crt
openssl pkcs12 -export \
-name kafka-broker \
-in broker.crt \
-inkey broker.key \
-certfile ca.crt \
-out /opt/kafka/secrets/kafka.keystore.p12 \
-passout "pass:$STORE_PASSWORD"
keytool -importcert -noprompt \
-alias kafka-ca \
-file ca.crt \
-keystore /opt/kafka/secrets/kafka.truststore.p12 \
-storetype PKCS12 \
-storepass "$STORE_PASSWORD"
正式域名必须同时出现在证书 SAN 和 KAFKA_HOST 中,否则客户端会报 hostname verification failed。ca.key 不应挂载进容器或分发给客户端;客户端只需要 CA 证书或只读 truststore。
创建容器内健康检查使用的管理员配置:
printf 'security.protocol=SASL_SSL\nsasl.mechanism=PLAIN\nsasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username="admin" password="%s";\nssl.truststore.location=/etc/kafka/secrets/kafka.truststore.p12\nssl.truststore.type=PKCS12\nssl.truststore.password=%s\n' \
"$ADMIN_PASSWORD" "$STORE_PASSWORD" \
> /opt/kafka/secrets/client-admin.properties
chmod 600 /opt/kafka/.env /opt/kafka/secrets/*
sudo chown -R 1000:1000 /opt/kafka/data /opt/kafka/secrets
创建 /opt/kafka/compose.yml:
services:
kafka:
image: apache/kafka:4.3.1
container_name: kafka
hostname: kafka
restart: unless-stopped
ports:
- "${KAFKA_BIND_IP:-127.0.0.1}:9094:9094"
volumes:
- ./data:/var/lib/kafka/data
- ./secrets:/etc/kafka/secrets:ro
environment:
CLUSTER_ID: ${CLUSTER_ID}
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,INTERNAL:SASL_SSL,EXTERNAL:SASL_SSL
KAFKA_LISTENERS: CONTROLLER://:29093,INTERNAL://:19094,EXTERNAL://:9094
KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka:19094,EXTERNAL://${KAFKA_HOST}:9094
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:29093
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL
KAFKA_SASL_ENABLED_MECHANISMS: PLAIN
KAFKA_SASL_MECHANISM_INTER_BROKER_PROTOCOL: PLAIN
KAFKA_OPTS: -Djava.security.auth.login.config=/etc/kafka/secrets/broker_jaas.conf
KAFKA_SSL_KEYSTORE_FILENAME: kafka.keystore.p12
KAFKA_SSL_KEYSTORE_TYPE: PKCS12
KAFKA_SSL_KEYSTORE_CREDENTIALS: kafka_keystore_creds
KAFKA_SSL_KEY_CREDENTIALS: kafka_ssl_key_creds
KAFKA_SSL_TRUSTSTORE_FILENAME: kafka.truststore.p12
KAFKA_SSL_TRUSTSTORE_TYPE: PKCS12
KAFKA_SSL_TRUSTSTORE_CREDENTIALS: kafka_truststore_creds
KAFKA_SSL_CLIENT_AUTH: none
KAFKA_AUTHORIZER_CLASS_NAME: org.apache.kafka.metadata.authorizer.StandardAuthorizer
KAFKA_SUPER_USERS: User:admin
KAFKA_ALLOW_EVERYONE_IF_NO_ACL_FOUND: "false"
KAFKA_LOG_DIRS: /var/lib/kafka/data
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
KAFKA_SHARE_COORDINATOR_STATE_TOPIC_REPLICATION_FACTOR: 1
KAFKA_SHARE_COORDINATOR_STATE_TOPIC_MIN_ISR: 1
KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
KAFKA_HEAP_OPTS: -Xms2g -Xmx2g
KAFKA_LOG_RETENTION_HOURS: 168
healthcheck:
test:
- CMD-SHELL
- /opt/kafka/bin/kafka-broker-api-versions.sh --bootstrap-server kafka:19094 --command-config /etc/kafka/secrets/client-admin.properties >/dev/null 2>&1
interval: 30s
timeout: 10s
retries: 10
start_period: 60s
mem_limit: 4g
cpus: 2.0
ulimits:
nofile:
soft: 100000
hard: 100000
logging:
driver: json-file
options:
max-size: 10m
max-file: "5"
controller 监听器只存在于容器网络,外部 9094 同时启用 SASL 与 TLS。KAFKA_ADVERTISED_LISTENERS 是客户端拿到的回连地址:域名、端口或 DNS 错一个,初次连接可能成功,随后仍会因为访问错误地址而超时。
启动前先做静默渲染,避免完整配置把密码写入终端日志:
cd /opt/kafka
docker compose config --quiet
docker compose pull
docker compose up -d
docker compose ps
docker compose logs --tail=100 kafka
先创建 3 分区、单副本的 orders Topic:
docker exec kafka /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server kafka:19094 \
--command-config /etc/kafka/secrets/client-admin.properties \
--create \
--topic orders \
--partitions 3 \
--replication-factor 1
为业务用户授予 Topic 读写和 Consumer Group 读取权限:
docker exec kafka /opt/kafka/bin/kafka-acls.sh \
--bootstrap-server kafka:19094 \
--command-config /etc/kafka/secrets/client-admin.properties \
--add --allow-principal User:app \
--operation Read --operation Write --operation Describe \
--topic orders
docker exec kafka /opt/kafka/bin/kafka-acls.sh \
--bootstrap-server kafka:19094 \
--command-config /etc/kafka/secrets/client-admin.properties \
--add --allow-principal User:app \
--operation Read \
--group orders-workers
复制 client-admin.properties 为业务客户端配置,把用户名改为 app、密码改为前面生成的 APP_PASSWORD。生产者和消费者都必须使用 security.protocol=SASL_SSL、sasl.mechanism=PLAIN 并信任 CA。不要用 ssl.endpoint.identification.algorithm= 关闭主机名校验来掩盖证书错误。
在容器内做一次最小验证:
printf 'order.created id=1001\n' | docker exec -i kafka \
/opt/kafka/bin/kafka-console-producer.sh \
--bootstrap-server kafka:19094 \
--producer.config /etc/kafka/secrets/client-admin.properties \
--topic orders
docker exec kafka /opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server kafka:19094 \
--consumer.config /etc/kafka/secrets/client-admin.properties \
--topic orders \
--from-beginning \
--max-messages 1
业务生产者应启用 acks=all;需要重试时还要理解幂等生产者、事务和消息键的语义,不能把“发送接口返回成功”直接等同于业务已完成。
本文把默认保留时间设为 168 小时。Kafka 会根据时间或空间策略删除旧日志,这不是永久归档。粗略估算原始数据量:
峰值写入 MB/s × 86400 × 保留天数 × 副本数
还要为索引、段文件、重放峰值和磁盘水位留余量。单节点磁盘使用率不要长期超过 70%–80%,否则压缩、日志删除滞后或消费者回放都可能迅速吃满剩余空间。
分区数决定并行度,但不是越多越好。先按消费者实例数和目标吞吐创建,再通过压测调整。分区只能增加,不能直接减少;如果依赖消息键保证顺序,增加分区还会改变键到分区的映射。
Kafka 默认不开放远程 JMX。生产环境可把 JMX Exporter 作为 Java Agent 接入,再由 Prometheus 抓取,但 JMX 或 metrics 端口仍应只在私网监听。至少监控:
UnderReplicatedPartitions和UnderMinIsrPartitionCount;- controller 是否唯一激活、controller 事件队列和选举异常;
- 请求延迟、错误率、BytesIn/BytesOut;
- Consumer Group 最大 lag 与消费速率;
- JVM GC、Heap、进程存活和文件描述符;
- 数据盘容量、inode、I/O 延迟与宿主机网络丢包。
单节点的副本指标不能替代可用性监控:只有一个 broker 时,即使 UnderReplicatedPartitions=0,主机宕机仍会让整个服务不可用。可结合 Prometheus + Grafana VPS 监控方案建立主机和服务告警。
Topic 副本用于可用性,不等于备份;同一集群中的误删、错误保留策略和应用写坏数据仍会传播。单节点可做一致性冷备份:
cd /opt/kafka
docker compose stop kafka
sudo tar -C /opt/kafka -czf \
"backups/kafka-data-$(date +%F-%H%M%S).tar.gz" \
data
docker compose start kafka
这是有停机窗口的快照,包含日志和 KRaft 元数据。恢复时应先在隔离目录验证,保持相同 Cluster ID、节点配置和文件权限,不要直接覆盖仍在运行的生产目录。
JAAS、CA、密钥库和客户端配置要单独做加密备份,并限制访问;把它们与数据压缩包明文放在同一 VPS 上没有灾备价值。要求低 RPO/RTO 时,应在第二套独立 Kafka 集群上使用 MirrorMaker 2 复制 Topic 和 Consumer Group checkpoint,并定期演练切换。
升级前先阅读目标版本的 release notes 和 upgrade guide,确认客户端、插件、Kafka Connect 与 broker 协议兼容。单节点没有滚动升级能力,步骤应包含维护窗口:
cd /opt/kafka
docker compose config --quiet
docker compose exec kafka /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server kafka:19094 \
--command-config /etc/kafka/secrets/client-admin.properties \
--list
docker compose pull
docker compose up -d
docker compose ps
docker compose logs --tail=100 kafka
不要直接使用 latest。先固定补丁版本、做冷备份并验证回滚条件;跨大版本升级还要按官方顺序处理 metadata version 与 KRaft feature version,不能只改镜像标签。
最常见原因是 advertised.listeners 返回了客户端无法解析或无法路由的地址。检查:
getent hosts mq.example.com
nc -vz mq.example.com 9094
域名必须从客户端解析到正确的公网或 WireGuard 私网 IP。再核对云安全组、主机防火墙和 Docker 端口绑定;可按 VPS 外网端口排查流程逐层检查。
客户端连接的主机名不在证书 SAN 中。用下面命令检查证书,不要关闭主机名验证:
openssl s_client -connect mq.example.com:9094 \
-servername mq.example.com \
-showcerts </dev/null
确认客户端机制为 PLAIN、协议为 SASL_SSL,并检查 JAAS 中是否存在对应 user_用户名。修改 broker_jaas.conf 后需要重启 broker;轮换密码时应先更新客户端,再移除旧凭据,避免集中断连。
通常是复用了旧数据目录,却改了 .env 中的 CLUSTER_ID,或把其他集群的数据盘挂了进来。停止服务,核对备份、meta.properties 和原 Cluster ID;不要用重新格式化数据目录来“修复”包含业务数据的节点。
检查宿主机目录所有者和只读 secrets:
cd /opt/kafka
docker compose logs --tail=200 kafka
namei -l /opt/kafka/data /opt/kafka/secrets
sudo chown -R 1000:1000 /opt/kafka/data /opt/kafka/secrets
不需要。Kafka 4.x 使用 KRaft 管理集群元数据。不要把旧版 ZooKeeper Compose 配置与本文混用。
可以用于开发和低吞吐验证,但 Heap 应控制在约 1GB,并严格限制保留时间和磁盘水位。小型生产更建议 8GB 内存和 NVMe,并通过真实消息压测决定容量。
只能用于能接受停机与数据恢复的小型业务。单节点没有 broker、controller 或主机级容错;关键业务应迁移到至少 3 个独立故障域。
不建议。优先通过 WireGuard 等私网连接,并使用 SASL_SSL、ACL、云安全组来源限制和主机防火墙。controller 端口绝不能开放给客户端网络。
只有搭配 TLS 才适合使用。PLAIN 本身是用户名密码机制,官方明确建议通过 SSL 传输;更复杂的团队可改用 SCRAM、OAuthBearer 或外部认证回调。
不等于。副本主要解决节点故障,无法防止误删、错误保留策略和逻辑损坏。还需要跨集群复制、冷备份或独立归档,并定期验证恢复。
