ELFK对接Kafka & KRaft模式集群
4146 字
21 分钟
ELFK对接Kafka & KRaft模式集群

ELFK对接Kafka && KRaft集群
[TOC]
环境规划
在上一篇笔记中,我们已经在 Elk01/02/03 上部署了 ZK 模式的 Kafka 集群 + ES7 HTTPS 集群 本篇承接上文,做三件事:
| 阶段 | 虚拟机 | 内容 |
|---|---|---|
| ① | Elk01/02/03 | ELFK对接Kafka:Filebeat→Kafka→Logstash→ES 全链路 |
| ② | Elk02 | Go生成100w条nginx日志 + ELFK 分析处理 |
| ③ | ESnew-01/02/03 | KRaft模式 Kafka 集群(无 Zookeeper) |
| 主机 | IP | 角色 |
|---|---|---|
| Elk01 | 10.0.0.6 | ES + ZK + Kafka(broker.id=1) |
| Elk02 | 10.0.0.7 | ES + ZK + Kafka(broker.id=2) + Filebeat + Go日志生成 |
| Elk03 | 10.0.0.8 | ES + ZK + Kafka(broker.id=3) + Logstash |
| ESnew-01 | 10.0.0.9 | KRaft Kafka(broker+controller) |
| ESnew-02 | 10.0.0.10 | KRaft Kafka(broker+controller) |
| ESnew-03 | 10.0.0.11 | KRaft Kafka(broker+controller) |
Kafka 在 ElasticStack 架构的位置

Note
📌 为什么加 Kafka?
- Filebeat 和 Logstash 之间插入 Kafka,削峰填谷 + 解耦
- Filebeat 只管往 Kafka 扔,Logstash 按自己节奏从 Kafka 消费
- 即使 Logstash 挂了,数据在 Kafka 里保留 7 天(默认),重启后从 offset 接着读
- 一个 topic 可以被多个消费者组独立消费
- 一组给 Logstash 写 ES,一组给 Spark 做分析
ElasticStack 对接 Kafka
启动 Kafka 集群
Elk01/02/03 上已安装 Kafka 3.9.2(ZK 模式),直接启动:
1)三台启动 Kafkaroot@Elk01 ~# source /etc/profile.d/kafka.sh && kafka-server-start.sh -daemon $KAFKA_HOME/config/server.properties# 三台都执行 ⚠️root@Elk01 ~# systemctl status zookeeper.service● zookeeper.service - Apache Zookeeper Loaded: loaded (/usr/lib/systemd/system/zookeeper.service; enabled; preset: enabled) Active: active (running) since Wed 2026-08-05 19:45:03 CST; 36min ago# 我这里是因为有zookeeper服务,开机自启动!
2)验证端口root@Elk01 ~# ss -lnt | grep -E "9092|2181|9200|9300"LISTEN 0 50 *:9092 *:*LISTEN 0 50 *:2181 *:*LISTEN 0 4096 *:9200 *:*LISTEN 0 4096 *:9300 *:*✅️ Kafka + ZK + ES 全部在线Filebeat 写入数据到 Kafka
1)创建 topic(5 分区 2 副本)root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic elfk-kafka --create --partitions 5 --replication-factor 2Created topic elfk-kafka.root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic elfk-kafka --describeTopic: elfk-kafka PartitionCount: 5 ReplicationFactor: 2 Configs: Topic: elfk-kafka Partition: 0 Leader: 3 Replicas: 3,2 Isr: 3,2 Topic: elfk-kafka Partition: 1 Leader: 1 Replicas: 1,3 Isr: 1,3 Topic: elfk-kafka Partition: 2 Leader: 2 Replicas: 2,1 Isr: 2,1 Topic: elfk-kafka Partition: 3 Leader: 3 Replicas: 3,1 Isr: 3,1 Topic: elfk-kafka Partition: 4 Leader: 1 Replicas: 1,2 Isr: 1,2
2)编写 Filebeat 配置(Elk02 上)root@Elk02 ~# vim /etc/filebeat/config/tcp-to-kafka.yamlfilebeat.inputs:- type: tcp host: "0.0.0.0:9000"
output.kafka: hosts: ["10.0.0.6:9092", "10.0.0.7:9092", "10.0.0.8:9092"] topic: elfk-kafka
3)启动 Filebeatroot@Elk02 ~# killall -9 filebeat 2>/dev/nullroot@Elk02 ~# nohup filebeat -e -c /etc/filebeat/config/tcp-to-kafka.yaml &> /tmp/fb-kafka.log &root@Elk02 ~# tail -1 /tmp/fb-kafka.logStarted listening for TCP connection {"address": "0.0.0.0:9000"}
4)发送测试数据(从 Elk01 发给 Elk02 的 9000 端口)root@Elk01 ~# echo www.kpyun.com | nc -w 1 10.0.0.7 9000
5)Kafka 验证数据已收到oot@Elk01 ~# kafka-console-consumer.sh --bootstrap-server 10.0.0.7:9092 --topic elfk-kafka --from-beginning | jq{ "@timestamp": "2026-08-05T12:29:58.386Z", "@metadata": { "beat": "filebeat", "type": "_doc", "version": "7.17.29" }, "ecs": { "version": "1.12.0" }, "host": { "name": "Elk02" }, "agent": { "name": "Elk02", "type": "filebeat", "version": "7.17.29", "hostname": "Elk02", "ephemeral_id": "a64f2ce4-80fe-434d-bfd8-96158b89d572", "id": "0faa14ff-9e2b-4b69-8840-5e8cae268f15" }, "message": "www.kpyun.com", "log": { "source": { "address": "10.0.0.6:35930" } }, "input": { "type": "tcp" }}Processed a total of 1 messages✅️ Filebeat → Kafka 链路通了Logstash 从 Kafka 消费数据写入 ES
1)创建 Logstash 专用 api-keyroot@Elk01 ~# curl -s -u elastic:passwd -k -X POST 'https://10.0.0.6:9200/_security/api_key' \ -H 'Content-Type: application/json' \ -d '{"name":"elfk-lg","role_descriptors":{"logstash_writer":{"cluster":["monitor","manage_index_templates","manage_ilm"],"index":[{"names":["kpyun-es-apikey-haha*"],"privileges":["create_index","create_doc","index","read"]}]}}}' | jq{ "id": "MQj50Z8BjME-gAjitirQ", "name": "elfk-lg", "api_key": "93WakTr4Ruisga8oGRwSBQ", "encoded": "TVFqNTBaOEJqTUUtZ0FqaXRpclE6OTNXYWtUcjRSdWlzZ2E4b0dSd1NCUQ=="}root@Elk01 ~# echo 'TVFqNTBaOEJqTUUtZ0FqaXRpclE6OTNXYWtUcjRSdWlzZ2E4b0dSd1NCUQ==' | base64 -d ;echoMQj50Z8BjME-gAjitirQ:93WakTr4Ruisga8oGRwSBQ
2)编写 Logstash 配置(Elk03 上)root@Elk03 ~# vim /etc/logstash/conf.d/kafka-to-es.confinput { kafka { bootstrap_servers => "10.0.0.6:9092,10.0.0.7:9092,10.0.0.8:9092" topics => "elfk-kafka" group_id => "linux-001" # 也可以通过更改组id重新获取数据 auto_offset_reset => "earliest" # topic 从头开始消费 }}
filter { json { source => "message" } mutate { remove_field => [ "@version","agent","log","ecs","tags","input" ] }}
output { elasticsearch { hosts => ["https://10.0.0.6:9200","https://10.0.0.7:9200","https://10.0.0.8:9200"] index => "kpyun-es-apikey-haha-01" api_key => "MQj50Z8BjME-gAjitirQ:93WakTr4Ruisga8oGRwSBQ" ssl => true ssl_certificate_verification => false }}
3)手动创建索引(我们关掉了自动创建)root@Elk01 ~# curl -s -u elastic:passwd -k -X PUT "https://10.0.0.6:9200/kpyun-es-apikey-haha-01" -H "Content-Type: application/json"{"acknowledged":true,"shards_acknowledged":true,"index":"kpyun-es-apikey-haha-2026-08-06"}
4)启动 Logstashroot@Elk03 ~# pkill -f logstash; nohup /usr/share/logstash/bin/logstash -rf /etc/logstash/conf.d/kafka-to-es.conf &> /tmp/ls-kafka.log &root@Elk03 ~# tail -f /tmp/ls-kafka.log[ERROR] 2026-08-06 09:53:16.981 [pool-7-thread-1] jvm - Unknown garbage collector name {:name=>"G1 Concurrent GC"}# 本质是 Logstash 7.x 的 Ruby 层日志和 Java 层日志是两套体系。jvm.rb 用 logger.error() 直接输出,它不受 log4j2 控制✅ 改源码 —— 把 jvm.rb 第 23 行 logger.error 改成 logger.debug:root@Elk03 ~# sed -i 's/logger.error(\"Unknown garbage collector name\"/logger.debug(\"Unknown garbage collector name\"/' \ /usr/share/logstash/logstash-core/lib/logstash/instrument/periodic_poller/jvm.rbroot@Elk03 ~# pkill -f logstash; nohup /usr/share/logstash/bin/logstash -rf /etc/logstash/conf.d/kafka-to-es.conf &> /tmp/ls-kafka.log &# 重新启动
5)发送测试数据验证全链路root@Elk01 ~# echo "kpyun-elfk-test" | nc -w 1 10.0.0.7 9000
6)ES 查询确认数据落地root@Elk01 ~# curl -s -u elastic:passwd -k "https://10.0.0.6:9200/kpyun-es-apikey-haha-01/_search?size=2" | jq '.hits.hits[1]._source'`拉了两条数据,查看第二条数据`{ "message": "kpyun-elfk-test", "host": { "name": "Elk02" }, "@timestamp": "2026-08-06T01:47:03.262Z"}✅️ Filebeat → Kafka → Logstash → ES 全链路贯通!Tip
💡 核心链路回顾:
nc → Filebeat(TCP:9000) → Kafka(elfk-kafka) → Logstash → ES(HTTPS + api-key)ELFK 架构分析 100w 条 nginx 日志
安装 Go 环境
1)下载 go1.25.0jiuzhao@Ubuntu ~$ wget -P /home/jiuzhao/下载 https://golang.google.cn/dl/go1.25.0.linux-amd64.tar.gz# 官方源被墙,用 google.cn 镜像
2)拷贝并安装到 Elk02jiuzhao@Ubuntu ~$ scp /home/jiuzhao/下载/go1.25.0.linux-amd64.tar.gz Elk02:/tmp/root@Elk02 ~# tar xf /tmp/go1.25.0.linux-amd64.tar.gz -C /usr/local/
3)配置环境变量root@Elk02 ~# mkdir -p /jiuzhao/gopathroot@Elk02 ~# vim /etc/profile.d/go.sh#!/bin/bashexport GOROOT=/usr/local/go/export GOPROXY=https://goproxy.cn,directexport GOPATH=/jiuzhao/gopathexport PATH=$PATH:$GOROOT/bin:$GOPATHroot@Elk02 ~# source /etc/profile.d/go.sh && go versiongo version go1.25.0 linux/amd64使用 Go 生成 100w 条 Nginx 日志
1)创建工作目录 & 初始化 go moduleroot@Elk02 ~# mkdir -p /jiuzhao/code/devops && cd /jiuzhao/code/devopsroot@Elk02 devops# source /etc/profile.d/go.sh && go mod init devopsgo: creating new go.mod: module devopsroot@Elk02 devops# lltotal 12drwxr-xr-x 2 root root 4096 Aug 6 10:29 ./drwxr-xr-x 3 root root 4096 Aug 6 10:29 ../-rw-r--r-- 1 root root 25 Aug 6 10:29 go.mod
2)编写 main.goroot@Elk02 devops# vim main.gopackage main
import ( "bufio" "fmt" "math/rand" "net" "os" "sync" "time")
const totalRecords = 1000000const goroutines = 10
var ( ipPools = []string{ "39.0.0.0/8", "42.0.0.0/8", "49.0.0.0/8", "58.0.0.0/8", "101.0.0.0/8", "103.0.0.0/8", "106.0.0.0/8", "110.0.0.0/8", "112.0.0.0/8", "113.0.0.0/8", "114.0.0.0/8", "115.0.0.0/8", "116.0.0.0/8", "117.0.0.0/8", "118.0.0.0/8", "119.0.0.0/8", "120.0.0.0/8", "121.0.0.0/8", "122.0.0.0/8", "123.0.0.0/8", "124.0.0.0/8", "125.0.0.0/8", "202.0.0.0/8", "203.0.0.0/8", "210.0.0.0/8", "211.0.0.0/8", "218.0.0.0/8", "219.0.0.0/8", "220.0.0.0/8", "221.0.0.0/8", "222.0.0.0/8", "223.0.0.0/8", }
methods = []string{"GET", "POST", "PUT", "DELETE", "HEAD", "PATCH"} uris = []string{"/api/v1/users", "/api/v1/orders", "/admin/login", "/favicon.ico", "/index.html", "/static/js/app.js", "/static/css/style.css", "/images/logo.png"} protocols = []string{"HTTP/1.0", "HTTP/1.1", "HTTP/2.0"} statusCodes = []string{"200", "201", "204", "301", "302", "400", "401", "403", "404", "500", "502", "503"} referers = []string{"-", "https://www.google.com", "https://www.baidu.com", "https://www.oldboyedu.com", "https://github.com"} userAgents = []string{ "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36", "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/119.0.0.0 Safari/537.36", "Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/118.0.0.0 Safari/537.36", "Mozilla/5.0 (iPhone; CPU iPhone OS 17_0 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) CriOS/120.0.6099.119 Mobile/15E148 Safari/604.1", "Mozilla/5.0 (Linux; Android 14; SM-G998B) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.6099.43 Mobile Safari/537.36", "Mozilla/5.0 (iPad; CPU OS 17_0 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.0 Mobile/15E148 Safari/604.1", "curl/7.68.0", "python-requests/2.28.0", })
func randomIP(ipNet *net.IPNet, rng *rand.Rand) string { ip := make(net.IP, len(ipNet.IP)) copy(ip, ipNet.IP) for i := range ip { ip[i] |= byte(rng.Intn(256)) &^ ipNet.Mask[i] } return ip.String()}
var mu sync.Mutex
func generateLogs(wg *sync.WaitGroup, writer *bufio.Writer, count int) { defer wg.Done()
startTime := time.Date(2023, 8, 5, 16, 47, 33, 0, time.FixedZone("CST", 8*3600)) endTime := time.Date(2026, 8, 5, 16, 47, 33, 0, time.FixedZone("CST", 8*3600)) duration := endTime.Unix() - startTime.Unix()
rng := rand.New(rand.NewSource(time.Now().UnixNano()))
for i := 0; i < count; i++ { mu.Lock() idx := rng.Intn(len(ipPools)) _, ipNet, _ := net.ParseCIDR(ipPools[idx]) remoteAddr := randomIP(ipNet, rng) remoteUser := "-" ts := time.Unix(startTime.Unix()+rng.Int63n(duration), 0) timeLocal := ts.Format("02/Jan/2006:15:04:05 -0700") method := methods[rng.Intn(len(methods))] uri := uris[rng.Intn(len(uris))] protocol := protocols[rng.Intn(len(protocols))] request := fmt.Sprintf("%s %s %s", method, uri, protocol) status := statusCodes[rng.Intn(len(statusCodes))] bodyBytes := rng.Intn(4096) + 64 referer := referers[rng.Intn(len(referers))] ua := userAgents[rng.Intn(len(userAgents))] mu.Unlock()
line := fmt.Sprintf(`%s %s - [%s] "%s" %s %d "%s" "%s"`+"\n", remoteAddr, remoteUser, timeLocal, request, status, bodyBytes, referer, ua)
mu.Lock() writer.WriteString(line) mu.Unlock() }}
func main() { fmt.Printf("开始生成%d条Nginx访问日志(时间范围:2023-08-05 16:47:33 至 2026-08-05 16:47:33)...\n", totalRecords)
file, err := os.Create("nginx_access.log") if err != nil { panic(err) } defer file.Close()
writer := bufio.NewWriterSize(file, 64*1024) defer writer.Flush()
start := time.Now() perGoroutine := totalRecords / goroutines
var wg sync.WaitGroup for g := 0; g < goroutines; g++ { wg.Add(1) go generateLogs(&wg, writer, perGoroutine) } wg.Wait() writer.Flush()
elapsed := time.Since(start)
stat, _ := os.Stat("nginx_access.log") fmt.Printf("日志生成完成!\n总记录数:%d条\n耗时:%v\n文件大小:%.2f MB\n日志文件路径:nginx_access.log\n", totalRecords, elapsed, float64(stat.Size())/1024/1024)}3)运行生成root@Elk02 devops# go run main.go# 近 3 年均匀分布开始生成1000000条Nginx访问日志(时间范围:2023-08-05 16:47:33 至 2026-08-05 16:47:33)...日志生成完成!总记录数:1000000条耗时:2.182946241s文件大小:196.01 MB日志文件路径:nginx_access.log
4)验证日志格式root@Elk02 devops# wc -l nginx_access.log1000000 nginx_access.logroot@Elk02 devops# tail -1 nginx_access.log211.98.210.150 - - [07/May/2025:17:13:31 +0800] "DELETE /admin/login HTTP/2.0" 201 409 "https://github.com" "Mozilla/5.0 (iPhone; CPU iPhone OS 17_0 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) CriOS/120.0.6099.119 Mobile/15E148 Safari/604.1"ELFK 架构分析处理
1)创建 topicroot@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic elfk-kafka-nginx --create --partitions 5 --replication-factor 2Created topic elfk-kafka-nginx.root@Elk02 devops# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic elfk-kafka-nginx --describeTopic: elfk-kafka-nginx PartitionCount: 5 ReplicationFactor: 2 Configs: Topic: elfk-kafka-nginx Partition: 0 Leader: 2 Replicas: 2,3 Isr: 2,3 Topic: elfk-kafka-nginx Partition: 1 Leader: 3 Replicas: 3,1 Isr: 3,1 Topic: elfk-kafka-nginx Partition: 2 Leader: 1 Replicas: 1,2 Isr: 1,2 Topic: elfk-kafka-nginx Partition: 3 Leader: 2 Replicas: 2,1 Isr: 2,1 Topic: elfk-kafka-nginx Partition: 4 Leader: 3 Replicas: 3,2 Isr: 3,2
2)Filebeat 配置(filestream 读取日志文件)root@Elk02 ~# vim /etc/filebeat/config/nginx-to-kafka.yamlfilebeat.inputs:- type: filestream paths: - /jiuzhao/code/devops/nginx_access.log
output.kafka: hosts: ["10.0.0.6:9092", "10.0.0.7:9092", "10.0.0.8:9092"] topic: elfk-kafka-nginx
root@Elk02 ~# killall -9 filebeatroot@Elk02 ~# nohup filebeat -e -c /etc/filebeat/config/nginx-to-kafka.yaml &> /tmp/fb-nginx.log &# 30 秒内 100w 条全部送入 Kafka ✅️
3)创建新的 api-key(给 nginx 索引写权限)root@Elk01 ~# curl -s -u elastic:passwd -k -X POST 'https://10.0.0.6:9200/_security/api_key' \ -H 'Content-Type: application/json' \ -d '{"name":"elfk-nginx","role_descriptors":{"logstash_writer":{"cluster":["monitor","manage_index_templates","manage_ilm"],"index":[{"names":["kpyun-kafka-nginx*"],"privileges":["create_index","create_doc","index","read"]}]}}}' | jq{ "id": "8sPK1Z8BlR5ZL1RhGdS4", "name": "elfk-nginx", "api_key": "yUDrmCLvQouUYTV24OYIkw", "encoded": "OHNQSzFaOEJsUjVaTDFSaEdkUzQ6eVVEcm1DTHZRb3VVWVRWMjRPWUlrdw=="}root@Elk01 ~# echo 'OHNQSzFaOEJsUjVaTDFSaEdkUzQ6eVVEcm1DTHZRb3VVWVRWMjRPWUlrdw==' | base64 -d ;echo8sPK1Z8BlR5ZL1RhGdS4:yUDrmCLvQouUYTV24OYIkw
4)Logstash 配置(grok + date + useragent + geoip)root@Elk03 ~# vim /etc/logstash/conf.d/kafka-to-es.confinput { kafka { bootstrap_servers => "10.0.0.6:9092,10.0.0.7:9092,10.0.0.8:9092" topics => "elfk-kafka-nginx" group_id => "linux-006" auto_offset_reset => "earliest" }}
filter { json { source => "message" } mutate { remove_field => [ "@version","agent","log","ecs","tags","input" ] } grok { match => { "message" => "%{HTTPD_COMBINEDLOG}" } # 基于正则套模板 } date { match => [ "timestamp", "dd/MMM/yyyy:HH:mm:ss Z" ] target => "@timestamp" } useragent { source => "agent" target => "kpyun-agent" } geoip { source => "clientip" database => "/root/GeoLite2-City_20250311/GeoLite2-City.mmdb" default_database_type => "City" }}
output {
# stdout {# codec => "rubydebug"# }
elasticsearch { hosts => ["https://10.0.0.6:9200","https://10.0.0.7:9200","https://10.0.0.8:9200"] index => "kpyun-kafka-nginx-test" api_key => "8sPK1Z8BlR5ZL1RhGdS4:yUDrmCLvQouUYTV24OYIkw" ssl => true ssl_certificate_verification => false }}| 组件 | 角色 | 说明 |
|---|---|---|
grok | nginx 日志解析 | %{HTTPD_COMBINEDLOG} 提取 clientip/timestamp/verb/request/httpversion…字段 |
date | 时间戳转换 | 把 02/Jan/2006:15:04:05 -0700 转为 @timestamp |
useragent | UA 分析 | 用 grok 提取的 agent 字段解析,输出浏览器/版本/OS/设备到 kpyun-agent |
geoip | IP 地理定位 | 用 grok 提取的 clientip 查 GeoLite2-City.mmdb,输出国家/城市/坐标 |
5)创建 索引模板(映射设置)+ ES 索引💡 可以通过 Kibana WebUI 创建'索引模板'索引模式:kpyun-kafka-nginx*索引设置: "number_of_shards": 7, # 主分片数 "number_of_replicas": 0 # 副本数映射设置: - `geoip.location` 为 `地理坐标点(geo_point)` - `bytes` 为 `数值 --> 整型(integer)` - `clientip` 为 `IP`=================================================# 先创建索引模板,再创建索引root@Elk01 ~# curl -s -u elastic:passwd -k -X PUT "https://10.0.0.6:9200/kpyun-kafka-nginx-test"

6)启动 Logstashroot@Elk03 ~# nohup /usr/share/logstash/bin/logstash -rf /etc/logstash/conf.d/kafka-to-es.conf &> /tmp/ls-nginx.log &
7)验证解析结果(100w 条全部写入)root@Elk03 ~# curl -s -u elastic:passwd -k "https://10.0.0.6:9200/kpyun-kafka-nginx-test/_count"{"count":1000000,"_shards":{"total":7,"successful":7,"skipped":0,"failed":0}}
8)查看样例文档root@Elk03 ~# curl -s -u elastic:passwd -k "https://10.0.0.6:9200/kpyun-kafka-nginx-test/_search?size=1" | jq '.hits.hits[0]._source'{ "timestamp": "10/Aug/2024:18:18:45 +0800", "httpversion": "1.0", "message": "210.175.8.194 - - [10/Aug/2024:18:18:45 +0800] \"GET /index.html HTTP/1.0\" 404 3927 \"https://www.baidu.com\" \"Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/119.0.0.0 Safari/537.36\"", "verb": "GET", "response": "404", "ident": "-", "geoip": { "continent_code": "AS", "country_code2": "JP", "timezone": "Asia/Tokyo", "latitude": 35.69, "longitude": 139.69, "country_code3": "JP", "location": { "lon": 139.69, "lat": 35.69 }, "country_name": "Japan", "ip": "210.175.8.194" }, "kpyun-agent": { "name": "Chrome", "major": "119", "version": "119.0.0.0", "minor": "0", "os": "Mac OS X", "os_version": "10.15.7", "os_name": "Mac OS X", "patch": "0", "os_minor": "15", "device": "Mac", "os_full": "Mac OS X 10.15.7", "os_major": "10", "os_patch": "7" }, "@timestamp": "2024-08-10T10:18:45.000Z", "host": { "name": "Elk02" }, "clientip": "210.175.8.194", "request": "/index.html", "referrer": "\"https://www.baidu.com\"", "bytes": "3927", "auth": "-", "agent": "\"Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/119.0.0.0 Safari/537.36\""}✅️ grok 解析正确、geoip 定位到日本、useragent 识别出 Chrome 119 + Mac OS X 10.15.7!Note
📌 解析效果总结
| 字段 | 示例值 | 来源 |
|---|---|---|
clientip | 210.175.8.194 | grok 从 nginx 日志提取 |
verb | GET | grok 提取 HTTP 方法 |
request | /index.html | grok 提取请求路径 |
response | 404 | grok 提取状态码 |
bytes | 3927(integer) | grok 提取 body 大小,索引模板映射为整型 |
httpversion | 1.0 | grok 提取 HTTP 协议版本 |
geoip.country_name | Japan | GeoLite2 数据库查 IP |
geoip.city_name | Nishiterao | GeoLite2 数据库查城市 |
geoip.location | {lat:35.69, lon:139.69} | geo_point 类型,可做地图可视化 |
kpyun-agent.name | Chrome | useragent 插件解析浏览器名称 |
kpyun-agent.os_full | Mac OS X 10.15.7 | useragent 插件解析完整 OS 信息 |
kpyun-agent.device | Mac | useragent 插件解析设备类型 |
Kibana出图展示





KRaft 模式部署 Kafka 集群
Zookeeper 和 Kafka 的关系

Note
ZooKeeper 做了什么?
- 集群元数据管理:topic、partition、ACL 等信息的持久化存储
- 控制器选举:多个 broker 谁当 controller
- 消费者组协调:消费者组的 offset 记录(老版本,新版本用了
__consumer_offsets内置 topic)
ZK 为 Kafka 提供了 可靠的分布式协调服务
为什么移除 Zookeeper
Important
📌 四宗罪
| 问题 | 说明 |
|---|---|
| 复杂性增加 | ZK 是独立组件,需单独部署和维护,运维两套分布式系统 |
| 性能瓶颈 | 高负载下 ZK 可能成为系统瓶颈,分区数增加 → 元数据监听延迟变大 |
| 一致性问题 | ZK 和 Kafka 的元数据同步不够高效,可能状态不一致 |
| 自立门户 | Kafka 生态壮大后不再依赖外部组件,减少被卡脖子风险 |
Tip
💡 版本演进时间线
- Kafka 2.8+:首次引入 KRaft 模式(预览 Preview)
- Kafka 3.x:ZK 和 KRaft 共存,生产主流仍是 ZK 模式
- Kafka 4.0+:彻底移除 ZK,自管理元数据信息,只支持 KRaft
KRaft 模式的核心优势
Important
📌 KRaft = 内置 ZK
KRaft 是 Kafka 内置的基于 Raft 一致性协议的元数据管理模块
- 👑 Controller 节点:管理集群元数据(替代 ZK),通过 Raft 协议做共识
- 📦 Broker 节点:消息照样落盘在 broker 本地,只是元数据从 ZK 搬到了 KRaft 内置的
@metadatatopic - 同一个节点可以同时是 broker + controller(复用)
| 对比维度 | ZK 模式 | KRaft 模式 |
|---|---|---|
| 外部依赖 | 需要 ZK 集群 | 不需要 |
| 元数据存储 | ZK | Kafka 内置 @metadata topic |
| 控制器选举 | ZK 协调 | Raft 协议自选举 |
| 故障转移 | 依赖 ZK 感知 | 更快(内置 Raft) |
| 扩展性 | 受 ZK 限制 | 数百万分区 |
| 部署复杂度 | 高(两套系统) | 低(一套系统) |
Tip
💡 Kafka 4.0 后的世界:部署 Kafka 就像部署 ES 一样——一个进程搞定全部,不再需要先装 ZK 再装 Kafka
KRaft 模式集群部署实战
环境:ESnew-01/02/03(从 init 快照恢复,干净系统,2C/2G)
1)恢复快照(ESnew 三台)jiuzhao@Ubuntu ~$ for i in ESnew-0{1..3} ;do virsh snapshot-revert $i init ;donejiuzhao@Ubuntu ~$ for i in ESnew-0{1..3} ;do virsh start $i ;done
2)安装 JDK8 & kafka(ESnew 上是干净系统,没有 JDK)jiuzhao@Ubuntu ~$ scp jdk8-temurin.tar.gz ESnew-01:/tmp/jiuzhao@Ubuntu ~$ scp kafka_2.13-3.9.2.tgz ESnew-01:/tmp/root@ESnew-01 ~# tar xzf /tmp/jdk8-temurin.tar.gz -C /usr/local/root@ESnew-01 ~# tar xf /tmp/kafka_2.13-3.9.2.tgz -C /usr/local/
3)配置环境变量root@ESnew-01 ~# vim /etc/profile.d/kafka.sh#!/bin/bashexport KAFKA_HOME=/usr/local/kafka_2.13-3.9.2export JAVA_HOME=/usr/local/jdk8u502-b07export PATH=$PATH:$KAFKA_HOME/bin:$JAVA_HOME/bin
4)预调 JVM 堆(2G 机器避免 OOM)root@ESnew-01 ~# sed -i 's/KAFKA_HEAP_OPTS="-Xmx1G -Xms1G"/KAFKA_HEAP_OPTS="-Xmx256m -Xms256m"/' \ $KAFKA_HOME/bin/kafka-server-start.sh
5)修改 server.properties — KRaft 核心配置root@ESnew-01 ~# source /etc/profile.d/kafka.shroot@ESnew-01 ~# vim $KAFKA_HOME/config/server.propertiesbroker.id=1 # ESnew-02: 2, ESnew-03: 3advertised.listeners=PLAINTEXT://10.0.0.9:9092
listeners=PLAINTEXT://:9092,CONTROLLER://:9093# 👆 两个端口:9092 给客户端,9093 给控制器通信
log.dirs=/var/lib/kafka
# ZK 配置必须注释掉!# zookeeper.connect=localhost:2181# zookeeper.connection.timeout.ms=18000
# KRaft 核心三项配置:(在末尾追加)process.roles=broker,controller# 节点角色(broker + controller 复用)controller.listener.names=CONTROLLER# controller 通信用的 listener 名称controller.quorum.voters=1@10.0.0.9:9093,2@10.0.0.10:9093,3@10.0.0.11:9093# 👆 quorum 投票者列表——KRaft 模式"不需要 ZK"的底气就在这里Important
📌 KRaft 三项核心配置解读
| 配置 | 值 | 说明 |
|---|---|---|
process.roles | broker,controller | 节点既是 broker(存数据)也是 controller(管元数据) |
controller.listener.names | CONTROLLER | controller 间通信用的 listener 名,必须存在 listeners 中 |
controller.quorum.voters | 1@10.0.0.9:9093,... | Raft 协议的投票者列表,格式 node.id@IP:CONTROLLER_PORT |
6)各节点创建数据目录 + 添加 hosts 解析root@ESnew-01 ~# mkdir -p /var/lib/kafkaroot@ESnew-01 ~# cat >> /etc/hosts <<EOF10.0.0.9 ESnew-0110.0.0.10 ESnew-0210.0.0.11 ESnew-03EOF# 三台都加 hosts,KRaft 需要主机名解析
7)拷贝jdk、kafka、环境变量、hostsroot@ESnew-01 ~# scp -r /usr/local/{jdk8u502-b07,kafka_2.13-3.9.2} ESnew-02:/usr/localroot@ESnew-01 ~# scp -r /usr/local/{jdk8u502-b07,kafka_2.13-3.9.2} ESnew-03:/usr/localroot@ESnew-01 ~# scp /etc/profile.d/kafka.sh ESnew-02:/etc/profile.d/root@ESnew-01 ~# scp /etc/profile.d/kafka.sh ESnew-03:/etc/profile.d/root@ESnew-01 ~# scp /etc/hosts ESnew-02:/etc/root@ESnew-01 ~# scp /etc/hosts ESnew-03:/etc/
8)更改主配置 & 创建数据目录root@ESnew-02 ~# source /etc/profile.d/kafka.sh && vim $KAFKA_HOME/config/server.propertiesbroker.id=2# ESnew-02: 2, ESnew-03: 3advertised.listeners=PLAINTEXT://10.0.0.10:9092....root@ESnew-03 ~# source /etc/profile.d/kafka.sh && vim $KAFKA_HOME/config/server.propertiesbroker.id=3# ESnew-02: 2, ESnew-03: 3advertised.listeners=PLAINTEXT://10.0.0.11:9092root@ESnew-02 ~# mkdir -p /var/lib/kafkaroot@ESnew-03 ~# mkdir -p /var/lib/kafka
9)生成集群 UUIDroot@ESnew-01 ~# source /etc/profile.d/kafka.shroot@ESnew-01 ~# kafka-storage.sh random-uuidxsUe7WA7SXWRxOaDM92JmQ
10)三台初始化(格式化元数据目录)root@ESnew-01 ~# kafka-storage.sh format -t xsUe7WA7SXWRxOaDM92JmQ -c $KAFKA_HOME/config/server.properties`把集群 UUID 和 node.id 写进本地,有了"身份信息"后再启动 Kafka`Formatting metadata directory /var/lib/kafka with metadata.version 3.9-IV0.# 三台同时执行同一个 UUID
root@ESnew-01 ~# cat /var/lib/kafka/meta.properties#Thu Aug 06 21:09:30 CST 2026cluster.id=xsUe7WA7SXWRxOaDM92JmQ # 同一集群的 UUID 统一version=1directory.id=qqnY7cCi518hXOBd4EfjcQ # 每个节点的唯一目录 IDnode.id=1 # 节点 ID# ESnew-02: node.id=2, ESnew-03: node.id=3Caution
⚠️ 三台必须用同一个 UUID 格式化,UUID 是集群的”身份证”
11)后台启动所有节点root@ESnew-01 ~# kafka-server-start.sh -daemon $KAFKA_HOME/config/server.properties# 三台都执行
12)验证端口(9092 客户端 + 9093 控制器)root@ESnew-01 ~# ss -lnt | egrep "9092|9093"LISTEN 0 50 *:9092 *:*LISTEN 0 50 *:9093 *:*✅️ 两个端口:9092 对外服务,9093 controller 通信验证 KRaft 集群
1)创建 topicroot@ESnew-01 ~# kafka-topics.sh --bootstrap-server 10.0.0.9:9092 --topic kpyun-linux --createCreated topic kpyun-linux.
root@ESnew-01 ~# kafka-topics.sh --bootstrap-server 10.0.0.9:9092 --topic kpyun-linux --describeTopic: kpyun-linux PartitionCount: 1 ReplicationFactor: 1 Configs: Topic: kpyun-linux Partition: 0 Leader: 2 Replicas: 2 Isr: 2✅️ KRaft 模式 topic 正常
2)生产者写数据root@ESnew-01 ~# echo "hello-kraft-cluster" | kafka-console-producer.sh --bootstrap-server 10.0.0.9:9092 --topic kpyun-linux
3)消费者读数据(跨节点消费)root@ESnew-02 ~# kafka-console-consumer.sh --bootstrap-server 10.0.0.10:9092 --topic kpyun-linux --from-beginninghello-kraft-cluster✅️ KRaft 集群跨节点消费正常,无需 ZK!Important
📌 KRaft vs ZK 模式体验对比
| 对比项 | ZK 模式 | KRaft 模式 |
|---|---|---|
| 启动顺序 | 先启 ZK → 再启 Kafka | 直接启 Kafka |
| 初始化 | 启动即用(ZK 已准备好) | 先写入集群元数据,再启动 |
| 端口 | 只开 9092 | 开 9092 + 9093 |
| 配置文件 | zk.connect 指向 ZK | controller.quorum.voters |
互联网企业架构使用 ELFK 对接 Kafka

Note
📌 真实企业案例
场景:Ucloud 7 台 Kafka 集群性能差 → 迁移到酒仙桥线下机房 5 台
踩坑记录:
| 问题 | 现象 | 解决方案 |
|---|---|---|
| 数据丢失 | 酒仙桥专线 2G 打满 | 升级 10G 带宽,HDFS 迁移限速 |
| 数据延迟 | 单个 Logstash 高峰期消费慢 | 新增 Logstash 到同消费者组(rebalance 分担) |
| ES 不显示 | 数据映射(mapping)冲突 | 提前定义索引 mapping |
课后练习
在 ES7 集群,新增一个 elk94 节点,并查看 ES 集群节点列表文章分享
如果这篇文章对你有帮助,欢迎分享给更多人!
相关文章智能推荐
1
Zookeeper优化与Kafka集群实战
ES集群ZK JVM 调优与 zkUI 图形化管理,Kafka 消息队列单点/集群部署、生产者消费者验证、常用术语与脚本、消费者组 rebalance、丢数据原理(ISR/LEO/HW)、JVM 与参数优化、Kafbat UI 图形化管理
2
Kibana数据可视化与ELFK架构实战
ES集群Kibana 出图展示,EFK 架构采集 Web 集群日志,Filebeat filestream 多行匹配,ELFK 三级管道处理自研 App 日志(mutate+date 清洗),Kibana 业务指标分析(PV/交易额/SVIP/用户行为 分层),EFK/ELFK 故障排查指南
3
Logstash 插件实战与 ES 集群安全加固
ES集群Logstash grok/date/useragent/geoip 插件实战,ELFK 出图故障排查与索引模板修复,Logstash 多实例、if多分支与 pipeline,ES 集群认证与 RBAC 权限隔离
4
ElasticStack开篇
ES集群从单点部署到集群,集群术语,DSL 语句实操,Kibana 可视化,Filebeat 日志采集全链路,EFK 架构实战
5
Go语言开篇
Go语言Go 1.26.6 环境安装与环境变量配置,对比系统安装与家目录安装两种方式,含 GOROOT/GOPATH 概念与 GOPROXY 配置















