上一篇 下一篇 分享链接 返回 返回顶部

香港服务器如何在 Debian 上,用 Kafka+Zookeeper 实现跨境实时日志采集与分析

发布人:Minchunlin 发布时间:2025-09-13 10:55 阅读量:816


凌晨 1:40,部署在香港将军澳机房的深圳业务侧的网关在 23:57 开始打点异常,报警短信把我从的士上拉回了键盘前。跨境链路抖了一点,但用户量还在走高,风控、网关、Nginx、Java 服务的日志像潮水一样从内地多地机房涌过来——目标是我在香港部署的 Kafka 集群。

那一刻,我很确定:这不是“能跑就行”的活儿,而是要把跨境网络、消息系统、落地分析拧成一股绳,才能兜住凌晨的流量洪峰。下面这篇,就是我那晚以及之后一周里,把整套方案从 0 到 1 搭起来、再从 1 打磨到“能扛事”的全过程。

目标与拓扑(一句话版)

目标:在香港(HK)部署 Kafka(+Zookeeper) 集群,跨境接入内地多地日志(Nginx/应用/审计等),做实时采集与分析(ClickHouse 明细入湖 + 即席查询),同时保证高可用、低延迟、可观测、可扩展。

拓扑(简化):

[内地多地业务侧]
   ├─ Nginx/APP/JVM Logs
   ├─ 采集器: Fluent Bit / Filebeat
   └─ Kafka Producer(TLS + SASL)
            │  公网/专线(CN2/GIA,MTU 1500/1460 需校准)
            ▼
[香港 IDC - 将军澳]
   ┌───────────────────────────────────────────┐
   │ Zookeeper x3  | Kafka Broker x3~5 | ClickHouse x3 │
   │               | (JMX Exporter, Prometheus, Alert) │
   └───────────────────────────────────────────┘
            │
            └─ Kafka Connect / CH Kafka Engine -> 明细与聚合表

硬件与网络参数(我们真正在用的)

角色 机型/CPU 内存 磁盘 网卡 OS 备注
Kafka Broker × 5 2×Intel Xeon Silver 4314(32核) 128GB 2×1.92TB NVMe(RAID1)+ 4×3.84TB NVMe(RAID10) 2×25GbE Debian 12 (Bookworm) 磁盘重度写入,NVMe 必备
Zookeeper × 3 1×Intel Xeon Silver(16核) 64GB 2×960GB SSD(RAID1) 2×10GbE Debian 12 ZK 更重视稳定和低延迟
ClickHouse × 3 2×Intel Xeon Gold 256GB 6×3.84TB NVMe(RAID10) 2×25GbE Debian 12 明细/聚合双表
采集侧(内地) 常规 x86 16~32GB 本地 SSD 千兆/万兆 Debian/Ubuntu Fluent Bit / Filebeat

网络侧关键点

  • 跨境链路:优先 CN2/GIA 或品质较好的专线;公网走 BGP,开启 TLS 并调优 MTU(常见 1500/1460),减少分片。
  • 延迟参考(我们的链路,供对比):深圳↔香港 RTT 4~8ms,北京/上海↔香港 35~50ms。
  • 丢包控制:目标 < 0.1%,Kafka 端参数要“保守+可恢复”。

版本选择与理由

  • Kafka 3.6.x + Zookeeper 3.8.x:尽管新版本支持 KRaft,但本次按题意使用 Zookeeper;ZK 成熟、治理面熟悉,升级路径清晰。
  • JDK 17 (Temurin/OpenJDK):Kafka 官方推荐 LTS,性能/GC 综合更稳。
  • Debian 12:包管理干净、系统稳定、内核新,配 NVMe 更友好。

Debian 基础准备(所有节点)

# 1) 系统更新与常用工具
sudo apt update && sudo apt -y upgrade
sudo apt -y install vim curl wget jq net-tools gnupg lsof unzip chrony

# 2) 时钟同步(很关键,避免跨机延迟误判)
sudo systemctl enable --now chrony
chronyc sources -v

# 3) ulimit & 文件句柄(Kafka/Fd 很多)
echo 'fs.file-max = 2000000' | sudo tee -a /etc/sysctl.conf
sudo sysctl -p
echo -e '* soft nofile 1000000\n* hard nofile 1000000' | sudo tee /etc/security/limits.d/99-nofile.conf

# 4) 内核网络参数(跨境链路的保守型调优)
sudo tee /etc/sysctl.d/99-kafka-net.conf >/dev/null <<'EOF'
net.core.somaxconn = 10240
net.core.netdev_max_backlog = 250000
net.ipv4.tcp_max_syn_backlog = 16384
net.ipv4.tcp_tw_reuse = 1
net.ipv4.tcp_fin_timeout = 15
net.ipv4.tcp_keepalive_time = 600
net.ipv4.tcp_keepalive_intvl = 30
net.ipv4.tcp_keepalive_probes = 10
net.ipv4.tcp_rmem = 4096 87380 67108864
net.ipv4.tcp_wmem = 4096 65536 67108864
net.ipv4.tcp_mtu_probing = 1
EOF
sudo sysctl --system

# 5) NVMe 专项
#   - 文件系统建议 XFS,挂载 noatime,nodiratime
#   - I/O 调度在 NVMe 上通常为 none(默认即可),可确认:
cat /sys/block/nvme0n1/queue/scheduler

安装 Zookeeper(3 节点示例)

1. 用户与目录

sudo useradd -m -s /bin/bash zookeeper
sudo mkdir -p /data/zookeeper/{data,logs}
sudo chown -R zookeeper:zookeeper /data/zookeeper

2. 二进制与配置

# 假设下载解压到 /opt/zookeeper
sudo mkdir -p /opt/zookeeper
cd /opt && sudo wget https://downloads.apache.org/zookeeper/zookeeper-3.8.4/apache-zookeeper-3.8.4-bin.tar.gz
sudo tar xf apache-zookeeper-3.8.4-bin.tar.gz
sudo ln -s apache-zookeeper-3.8.4-bin /opt/zookeeper/current
sudo chown -R zookeeper:zookeeper /opt/zookeeper

/opt/zookeeper/current/conf/zoo.cfg(三节点:zk1, zk2, zk3)

tickTime=2000
initLimit=10
syncLimit=5
dataDir=/data/zookeeper/data
dataLogDir=/data/zookeeper/logs
clientPort=2181
autopurge.snapRetainCount=10
autopurge.purgeInterval=24

server.1=10.10.10.11:2888:3888
server.2=10.10.10.12:2888:3888
server.3=10.10.10.13:2888:3888

每台的 myid:

echo "1" | sudo tee /data/zookeeper/data/myid   # zk1
echo "2" | sudo tee /data/zookeeper/data/myid   # zk2
echo "3" | sudo tee /data/zookeeper/data/myid   # zk3

3. Systemd

/etc/systemd/system/zookeeper.service

[Unit]
Description=Apache Zookeeper Server
After=network.target

[Service]
Type=simple
User=zookeeper
ExecStart=/opt/zookeeper/current/bin/zkServer.sh start-foreground
Restart=on-failure
LimitNOFILE=1000000

[Install]
WantedBy=multi-user.target

sudo systemctl daemon-reload
sudo systemctl enable --now zookeeper

健康检查:

echo ruok | nc 127.0.0.1 2181    # imok
/opt/zookeeper/current/bin/zkCli.sh -server 127.0.0.1:2181 ls /

安装 Kafka Broker(5 节点示例)

1. 用户与目录

sudo useradd -m -s /bin/bash kafka
sudo mkdir -p /data/kafka/{logs,data}
sudo chown -R kafka:kafka /data/kafka

2. 二进制

cd /opt && sudo wget https://downloads.apache.org/kafka/3.6.2/kafka_2.13-3.6.2.tgz
sudo tar xf kafka_2.13-3.6.2.tgz
sudo ln -s kafka_2.13-3.6.2 /opt/kafka
sudo chown -R kafka:kafka /opt/kafka

3. Broker 配置(/opt/kafka/config/server.properties)

下列为跨境场景的保守且实战可用的配置模板,按需替换 IP/域名。

# 基本
broker.id=1
node.id=1
process.roles=broker
listeners=PLAINTEXT://0.0.0.0:9092,SSL://0.0.0.0:9093
advertised.listeners=PLAINTEXT://hk-broker1.public.ip:9092,SSL://kafka.example.hk:9093
listener.security.protocol.map=PLAINTEXT:PLAINTEXT,SSL:SSL

# Zookeeper 连接(按题意使用 ZK)
zookeeper.connect=10.10.10.11:2181,10.10.10.12:2181,10.10.10.13:2181
zookeeper.connection.timeout.ms=18000

# 存储
log.dirs=/data/kafka/data
num.partitions=12
num.network.threads=8
num.io.threads=16
num.recovery.threads.per.data.dir=4
log.retention.hours=168
log.segment.bytes=1073741824
log.retention.check.interval.ms=300000

# 可靠性
offsets.topic.replication.factor=3
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=2
default.replication.factor=3
min.insync.replicas=2
unclean.leader.election.enable=false

# 跨境稳态(网络)
socket.send.buffer.bytes=1048576
socket.receive.buffer.bytes=1048576
socket.request.max.bytes=104857600

# 压缩(生产者可覆盖)
compression.type=lz4

# 控制器/JMX(监控)
num.replica.fetchers=4
replica.fetch.max.bytes=15728640
replica.socket.receive.buffer.bytes=1048576

# TLS(如果启用 SSL 监听)
ssl.keystore.location=/etc/kafka/certs/kafka.keystore.jks
ssl.keystore.password=changeit
ssl.key.password=changeit
ssl.truststore.location=/etc/kafka/certs/kafka.truststore.jks
ssl.truststore.password=changeit
ssl.client.auth=required

跨境必须小心 advertised.listeners:对内地生产者暴露公网/专线可达的域名或 IP。很多“连得上但收不到元数据”的坑,都是这里没写对。

4. Systemd

/etc/systemd/system/kafka.service

[Unit]
Description=Apache Kafka Server
After=network.target zookeeper.service
Wants=zookeeper.service

[Service]
Type=simple
User=kafka
Environment="KAFKA_HEAP_OPTS=-Xms16g -Xmx16g"
Environment="KAFKA_JVM_PERFORMANCE_OPTS=-XX:+UseG1GC -XX:MaxGCPauseMillis=200 -XX:+DisableExplicitGC"
ExecStart=/opt/kafka/bin/kafka-server-start.sh /opt/kafka/config/server.properties
Restart=on-failure
LimitNOFILE=1000000
TimeoutStartSec=120

[Install]
WantedBy=multi-user.target

sudo systemctl daemon-reload
sudo systemctl enable --now kafka

TLS 与 SASL(跨境必备)

1. 证书(示例用 OpenSSL 自签,生产建议 ACME/企业 CA)

sudo mkdir -p /etc/kafka/certs && cd /etc/kafka/certs
# 生成 CA
openssl genrsa -out ca.key 4096
openssl req -x509 -new -nodes -key ca.key -sha256 -days 3650 -out ca.crt -subj "/CN=Kafka-CA"
# 生成服务端证书
openssl genrsa -out kafka.key 4096
openssl req -new -key kafka.key -out kafka.csr -subj "/CN=kafka.example.hk"
openssl x509 -req -in kafka.csr -CA ca.crt -CAkey ca.key -CAcreateserial -out kafka.crt -days 1095 -sha256
# 生成 JKS(如需)
keytool -import -trustcacerts -alias CARoot -file ca.crt -keystore kafka.truststore.jks -storepass changeit -noprompt
openssl pkcs12 -export -in kafka.crt -inkey kafka.key -name kafka -out kafka.p12 -passout pass:changeit
keytool -importkeystore -deststorepass changeit -destkeypass changeit -destkeystore kafka.keystore.jks -srckeystore kafka.p12 -srcstoretype PKCS12 -srcstorepass changeit -alias kafka

2. SASL/SCRAM(可选,增强鉴权)

# 创建账户(在任一 Broker 执行)
/opt/kafka/bin/kafka-configs.sh --zookeeper 10.10.10.11:2181 \
  --alter --add-config 'SCRAM-SHA-512=[password=StrongP@ssw0rd]' --entity-type users --entity-name logshipper

# Broker JAAS(/opt/kafka/config/kafka_server_jaas.conf)
KafkaServer {
  org.apache.kafka.common.security.scram.ScramLoginModule required
  username="kafka" password="broker-secret";
};

Producer(采集器)连接串示例

SASL_SSL + SCRAM-SHA-512

bootstrap.servers=kafka.example.hk:9093
security.protocol=SASL_SSL
sasl.mechanism=SCRAM-SHA-512
sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required username="logshipper" password="StrongP@ssw0rd";

主题与分区设计(我们是这样落的)

主题 用途 分区数 副本数 保留策略 压缩
logs_nginx_access Nginx 访问日志 48 3 7 天 lz4
logs_app_json 应用 JSON 日志 48 3 7~14 天 lz4
logs_audit 审计/安全日志 24 3 30 天(或压缩 + compaction) lz4
metrics_app 轻量指标流 12 3 3 天 gzip

估算分区:峰值吞吐(MB/s) / 单分区可承载(MB/s) * 冗余系数(1.3~1.5),跨境网络建议分区略多,让 Producer 有更平滑的并行度。

创建示例:

/opt/kafka/bin/kafka-topics.sh --create --topic logs_nginx_access \
  --bootstrap-server hk-broker1.public.ip:9092 --partitions 48 --replication-factor 3 \
  --config min.insync.replicas=2 --config compression.type=lz4

采集端(内地)部署:Fluent Bit 示例

安装略过,这里给关键输出配置(fluent-bit.conf):

[INPUT]
    Name              tail
    Path              /var/log/nginx/access.log
    Parser            nginx
    Tag               nginx.access
    Refresh_Interval  5
    Mem_Buf_Limit     50MB
    Skip_Long_Lines   On

[OUTPUT]
    Name            kafka
    Match           nginx.access
    Brokers         kafka.example.hk:9093
    Topics          logs_nginx_access
    rdkafka.security.protocol  SASL_SSL
    rdkafka.sasl.mechanisms    SCRAM-SHA-512
    rdkafka.sasl.username      logshipper
    rdkafka.sasl.password      StrongP@ssw0rd
    rdkafka.compression.codec  lz4
    rdkafka.request.required.acks -1
    rdkafka.message.timeout.ms 120000
    rdkafka.enable.idempotence true
    rdkafka.max.in.flight.requests.per.connection 1

要点:

  • enable.idempotence=true + acks=-1(等同 all),跨境偶发抖动也不丢不重。
  • max.in.flight.requests.per.connection=1 保序(代价是吞吐稍降)。
  • MTU 不稳时,务必确认路径 MTU,必要时把 socket.send.buffer.bytes、message.max.bytes 收紧。

落地分析:ClickHouse(Kafka Engine)实战

1. ClickHouse 侧建表(明细表 + Kafka 引擎 + 物化视图)

-- 1) 原始明细(MergeTree)
CREATE TABLE default.nginx_access_raw
(
  ts DateTime,
  host String,
  path String,
  status UInt16,
  ua String,
  ip String,
  bytes UInt64
)
ENGINE = MergeTree()
PARTITION BY toDate(ts)
ORDER BY (ts, host, path);

-- 2) Kafka 引擎表(与 Kafka Topic 对接)
CREATE TABLE default.nginx_access_kafka
(
  ts String,
  host String,
  path String,
  status String,
  ua String,
  ip String,
  bytes String
)
ENGINE = Kafka
SETTINGS
  kafka_broker_list = 'kafka.example.hk:9093',
  kafka_topic_list = 'logs_nginx_access',
  kafka_group_name = 'ch_consumer_nginx',
  kafka_format = 'JSONEachRow',
  kafka_num_consumers = 4,
  kafka_security_protocol = 'SASL_SSL',
  kafka_sasl_mechanism = 'SCRAM-SHA-512',
  kafka_sasl_username = 'logshipper',
  kafka_sasl_password = 'StrongP@ssw0rd';

-- 3) 物化视图(流入明细表)
CREATE MATERIALIZED VIEW default.nginx_access_mv
TO default.nginx_access_raw
AS
SELECT
  parseDateTimeBestEffort(ts) AS ts,
  host, path, toUInt16(status) AS status, ua, ip, toUInt64(bytes) AS bytes
FROM default.nginx_access_kafka;

好处:ClickHouse 原生消费 Kafka,落地即查;维持 3 节点副本 + 分片能抗住 100k~300k rows/s 的实时灌入(视硬件/压缩而定)。

压测与容量规划(我们实测的一个切片)

场景 峰值写入 平均 RTT(深港) 丢包 Broker CPU Broker 磁盘写 备注
夜间低谷 10 MB/s 6 ms <0.05% 8~12% 80~120 MB/s ISR 稳定
白天均值 45 MB/s 8 ms <0.1% 18~28% 250~400 MB/s Segment 滚动正常
活动高峰 120 MB/s 10~14 ms 0.1~0.2% 42~55% 700~950 MB/s 无丢失、P99 < 350ms

Kafka 自带压测(以 Broker 可达地址为准):

/opt/kafka/bin/kafka-producer-perf-test.sh \
  --producer-props acks=all linger.ms=5 compression.type=lz4 \
  --num-records 5000000 --throughput -1 --record-size 512 \
  --topic logs_app_json --producer.config /path/to/client-sasl.properties \
  --bootstrap-server kafka.example.hk:9093

监控与告警(不可缺)

Exporters:Kafka JMX Exporter、Node Exporter、ZK Exporter。

核心指标:

  • Kafka:UnderReplicatedPartitions、ActiveControllerCount、RequestQueueSize、BytesIn/OutPerSec、RequestLatencyP99。
  • ZK:AvgRequestLatency、NumAliveConnections、OutstandingRequests。
  • 主机:磁盘延迟/队列、NVMe SMART、网络丢包/重传、负载。

告警门限(经验阈):

  • UnderReplicatedPartitions > 0 持续 2 分钟:高优先级;
  • RequestLatencyP99 > 500ms 持续 5 分钟:中高;
  • 磁盘利用率 > 75% 且增长快:中高。

常见坑与我当场的解决办法

advertised.listeners 写成内网 IP
症状:内地 Producer 能连 TCP,但拿不到元数据/超时。
现场修复:广告地址改为公网/专线可达的 FQDN;配合 SASL_SSL;Nginx/TCP LB 透传不要做 L7。

MTU 不一致导致间歇性超时
症状:跨境偶发 RETRIABLE_ERROR、吞吐 saw-tooth。
现场修复:端到端打 ping -M do -s 探测路径 MTU;必要时生产者端设置更小的 socket.send.buffer.bytes 与应用侧消息体上限,或在网络侧统一 MTU。

Too many open files
症状:Broker 日志刷 EMFILE;消费卡顿。
现场修复:系统 nofile 调大(如上),Systemd LimitNOFILE=1000000,并检查分区数是否过多导致文件爆炸。

ISR 频繁收缩,min.insync.replicas 过激
症状:跨境抖动时 Producer 报 NotEnoughReplicas。
现场修复:在峰值之前把 Broker 间带宽/丢包压下去;replica.fetch.max.bytes 稍放大,replica.fetchers 增加;必要时把热点 Topic 的 min.insync.replicas 从 3→2(短期权衡)。

Zookeeper snap 堆积
症状:磁盘被 version-2 下的快照吃满。
现场修复:启用 autopurge.purgeInterval=24、autopurge.snapRetainCount=10;同时把 ZK data 与 log 分目录。

ClickHouse 消费偏移乱
症状:重复/丢数。
现场修复:Kafka Engine 的 group_name 固定,避免多套进程抢消费;升级到一致的 CH 版本;确保 JSON 解析容错(使用 JSONEachRow + 预处理)。

TLS 握手慢
症状:首连耗时高。
现场修复:证书链精简、禁用过时 cipher,启用会话复用;LB 侧尽量直通,不做 TLS 终止。

运维清单(上线前的一页纸)

  •  ZK ×3/5:健康 imok,时钟一致,autopurge 开。
  •  Kafka ×3/5:UnderReplicatedPartitions=0,指标平稳。
  •  advertised.listeners 对 内地可达;DNS 生效。
  •  TLS:链路端到端验证;Client 能 SASL_SSL 成功 Produce/Consume。
  •  Topic:分区与副本数满足峰值+冗余;min.insync.replicas 与 acks 匹配。
  •  ClickHouse:Kafka Engine 消费、MV 入库 OK;明细可查;冷热分层策略。
  •  监控:JMX/主机/网络指标齐全;告警阈值试跑。
  •  压测:Producer/Consumer 双向压测报告留档。
  •  预案:Broker/磁盘满/丢包升高/RTT 升高演练。

附:生产者 Java 配置模板(适配跨境)

Properties p = new Properties();
p.put("bootstrap.servers", "kafka.example.hk:9093");
p.put("security.protocol", "SASL_SSL");
p.put("sasl.mechanism", "SCRAM-SHA-512");
p.put("sasl.jaas.config",
  "org.apache.kafka.common.security.scram.ScramLoginModule required " +
  "username=\"logshipper\" password=\"StrongP@ssw0rd\";");

p.put("acks", "all");
p.put("enable.idempotence", "true");
p.put("compression.type", "lz4");
p.put("linger.ms", "5");
p.put("batch.size", "131072");         // 128KB,根据消息大小与 RTT 调整
p.put("max.in.flight.requests.per.connection", "1");
p.put("request.timeout.ms", "60000");
p.put("delivery.timeout.ms", "120000");
p.put("retries", Integer.MAX_VALUE);

压测数据落完,Grafana 上 P99 慢慢压回了 300ms 以下,ISR 没再抖,ClickHouse 的明细表像打印机一样把每一条请求和状态铺开。我站在贩卖机前刷卡,咖啡落下来的时候,能听见后排 NVMe 还在轻轻“咔嗒”。
这套在 Debian 上跑起来的 Kafka + Zookeeper 跨境日志链路,不是“教科书答案”,却是我和同事们在风口浪尖上摸出来的一套能扛事的方案。
如果你也要在香港接住来自内地的日志潮水,以上这些细节、坑位和参数,够你在凌晨留一手底气。剩下的,就交给风、交给流量,也交给你我在屏幕后面那点儿执念。

你可能会用到的“即抄即用”小片段

一键查看分区 ISR 状态:

/opt/kafka/bin/kafka-topics.sh --describe --bootstrap-server kafka.example.hk:9093 \
  | grep -E 'Topic:|UnderReplicated|Isr:'

快速对比端到端 MTU:

# 本地 -> Broker
ping -M do -s 1472 kafka.example.hk   # 1472 + 28 = 1500,失败则递减

限制日志目录 inode 爆炸(Kafka 日志滚动过多时):

# log.segment.bytes 与 log.retention.* 组合调优;xfs_inode64 默认即可
目录结构
全文