ELFK对接Kafka & KRaft模式集群

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

ELFK对接Kafka && KRaft集群#

[TOC]


环境规划#

在上一篇笔记中,我们已经在 Elk01/02/03 上部署了 ZK 模式的 Kafka 集群 + ES7 HTTPS 集群 本篇承接上文,做三件事:

阶段虚拟机内容
①Elk01/02/03ELFK对接Kafka:Filebeat→Kafka→Logstash→ES 全链路
②Elk02Go生成100w条nginx日志 + ELFK 分析处理
③ESnew-01/02/03KRaft模式 Kafka 集群(无 Zookeeper)
主机IP角色
Elk0110.0.0.6ES + ZK + Kafka(broker.id=1)
Elk0210.0.0.7ES + ZK + Kafka(broker.id=2) + Filebeat + Go日志生成
Elk0310.0.0.8ES + ZK + Kafka(broker.id=3) + Logstash
ESnew-0110.0.0.9KRaft Kafka(broker+controller)
ESnew-0210.0.0.10KRaft Kafka(broker+controller)
ESnew-0310.0.0.11KRaft Kafka(broker+controller)

Kafka 在 ElasticStack 架构的位置#

ElasticStack架构升级及MQ对比
ElasticStack架构升级及MQ对比

Filebeat

TCP:9000

Kafka

Topic: elfk-kafka

Logstash

ES7 集群

HTTPS + api-key

Filebeat

TCP:9000

Kafka

Topic: elfk-kafka

Logstash

ES7 集群

HTTPS + api-key

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 模式),直接启动:

Terminal window
1)三台启动 Kafka
root@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#

Terminal window
1)创建 topic(5 分区 2 副本)
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic elfk-kafka --create --partitions 5 --replication-factor 2
Created topic elfk-kafka.
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic elfk-kafka --describe
Topic: 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.yaml
filebeat.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)启动 Filebeat
root@Elk02 ~# killall -9 filebeat 2>/dev/null
root@Elk02 ~# nohup filebeat -e -c /etc/filebeat/config/tcp-to-kafka.yaml &> /tmp/fb-kafka.log &
root@Elk02 ~# tail -1 /tmp/fb-kafka.log
Started 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#

Terminal window
1)创建 Logstash 专用 api-key
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-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 ;echo
MQj50Z8BjME-gAjitirQ:93WakTr4Ruisga8oGRwSBQ
2)编写 Logstash 配置(Elk03 上)
root@Elk03 ~# vim /etc/logstash/conf.d/kafka-to-es.conf
input {
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)启动 Logstash
root@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.rb
root@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

💡 核心链路回顾:

Terminal window
nc → Filebeat(TCP:9000) → Kafka(elfk-kafka) → Logstash → ES(HTTPS + api-key)

ELFK 架构分析 100w 条 nginx 日志#

安装 Go 环境#

Terminal window
1)下载 go1.25.0
jiuzhao@Ubuntu ~$ wget -P /home/jiuzhao/下载 https://golang.google.cn/dl/go1.25.0.linux-amd64.tar.gz
# 官方源被墙,用 google.cn 镜像
2)拷贝并安装到 Elk02
jiuzhao@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/gopath
root@Elk02 ~# vim /etc/profile.d/go.sh
#!/bin/bash
export GOROOT=/usr/local/go/
export GOPROXY=https://goproxy.cn,direct
export GOPATH=/jiuzhao/gopath
export PATH=$PATH:$GOROOT/bin:$GOPATH
root@Elk02 ~# source /etc/profile.d/go.sh && go version
go version go1.25.0 linux/amd64

使用 Go 生成 100w 条 Nginx 日志#

Terminal window
1)创建工作目录 & 初始化 go module
root@Elk02 ~# mkdir -p /jiuzhao/code/devops && cd /jiuzhao/code/devops
root@Elk02 devops# source /etc/profile.d/go.sh && go mod init devops
go: creating new go.mod: module devops
root@Elk02 devops# ll
total 12
drwxr-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.go
root@Elk02 devops# vim main.go
package main
import (
"bufio"
"fmt"
"math/rand"
"net"
"os"
"sync"
"time"
)
const totalRecords = 1000000
const 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)
}
Terminal window
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.log
1000000 nginx_access.log
root@Elk02 devops# tail -1 nginx_access.log
211.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 架构分析处理#

Terminal window
1)创建 topic
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic elfk-kafka-nginx --create --partitions 5 --replication-factor 2
Created topic elfk-kafka-nginx.
root@Elk02 devops# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic elfk-kafka-nginx --describe
Topic: 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.yaml
filebeat.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 filebeat
root@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 ;echo
8sPK1Z8BlR5ZL1RhGdS4:yUDrmCLvQouUYTV24OYIkw
4)Logstash 配置(grok + date + useragent + geoip)
root@Elk03 ~# vim /etc/logstash/conf.d/kafka-to-es.conf
input {
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
}
}
组件角色说明
groknginx 日志解析%{HTTPD_COMBINEDLOG} 提取 clientip/timestamp/verb/request/httpversion…字段
date时间戳转换把 02/Jan/2006:15:04:05 -0700 转为 @timestamp
useragentUA 分析用 grok 提取的 agent 字段解析,输出浏览器/版本/OS/设备到 kpyun-agent
geoipIP 地理定位用 grok 提取的 clientip 查 GeoLite2-City.mmdb,输出国家/城市/坐标
Terminal window
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"

模板-映射设置
模板-映射设置

创建的索引
创建的索引

Terminal window
6)启动 Logstash
root@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

📌 解析效果总结

字段示例值来源
clientip210.175.8.194grok 从 nginx 日志提取
verbGETgrok 提取 HTTP 方法
request/index.htmlgrok 提取请求路径
response404grok 提取状态码
bytes3927(integer)grok 提取 body 大小,索引模板映射为整型
httpversion1.0grok 提取 HTTP 协议版本
geoip.country_nameJapanGeoLite2 数据库查 IP
geoip.city_nameNishiteraoGeoLite2 数据库查城市
geoip.location{lat:35.69, lon:139.69}geo_point 类型,可做地图可视化
kpyun-agent.nameChromeuseragent 插件解析浏览器名称
kpyun-agent.os_fullMac OS X 10.15.7useragent 插件解析完整 OS 信息
kpyun-agent.deviceMacuseragent 插件解析设备类型

Kibana出图展示#

带宽
带宽

操作系统
操作系统

访问国家
访问国家

访问设备
访问设备

地理地图
地理地图

KRaft 模式部署 Kafka 集群#

Zookeeper 和 Kafka 的关系#

zookeeper和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 模式的核心优势#

Kafka KRaft 模式

Broker + Controller 1

Broker + Controller 2

Broker + Controller 3

Kafka + ZK 模式

Broker 1

Zookeeper 1

Broker 2

Zookeeper 2

Broker 3

Zookeeper 3

Kafka KRaft 模式

Broker + Controller 1

Broker + Controller 2

Broker + Controller 3

Kafka + ZK 模式

Broker 1

Zookeeper 1

Broker 2

Zookeeper 2

Broker 3

Zookeeper 3

Important

📌 KRaft = 内置 ZK

KRaft 是 Kafka 内置的基于 Raft 一致性协议的元数据管理模块

  • 👑 Controller 节点:管理集群元数据(替代 ZK),通过 Raft 协议做共识
  • 📦 Broker 节点:消息照样落盘在 broker 本地,只是元数据从 ZK 搬到了 KRaft 内置的 @metadata topic
  • 同一个节点可以同时是 broker + controller(复用)
对比维度ZK 模式KRaft 模式
外部依赖需要 ZK 集群不需要
元数据存储ZKKafka 内置 @metadata topic
控制器选举ZK 协调Raft 协议自选举
故障转移依赖 ZK 感知更快(内置 Raft)
扩展性受 ZK 限制数百万分区
部署复杂度高(两套系统)低(一套系统)
Tip

💡 Kafka 4.0 后的世界:部署 Kafka 就像部署 ES 一样——一个进程搞定全部,不再需要先装 ZK 再装 Kafka

KRaft 模式集群部署实战#

环境:ESnew-01/02/03(从 init 快照恢复,干净系统,2C/2G)

Terminal window
1)恢复快照(ESnew 三台)
jiuzhao@Ubuntu ~$ for i in ESnew-0{1..3} ;do virsh snapshot-revert $i init ;done
jiuzhao@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/bash
export KAFKA_HOME=/usr/local/kafka_2.13-3.9.2
export JAVA_HOME=/usr/local/jdk8u502-b07
export 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.sh
root@ESnew-01 ~# vim $KAFKA_HOME/config/server.properties
broker.id=1 # ESnew-02: 2, ESnew-03: 3
advertised.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.rolesbroker,controller节点既是 broker(存数据)也是 controller(管元数据)
controller.listener.namesCONTROLLERcontroller 间通信用的 listener 名,必须存在 listeners 中
controller.quorum.voters1@10.0.0.9:9093,...Raft 协议的投票者列表,格式 node.id@IP:CONTROLLER_PORT
Terminal window
6)各节点创建数据目录 + 添加 hosts 解析
root@ESnew-01 ~# mkdir -p /var/lib/kafka
root@ESnew-01 ~# cat >> /etc/hosts <<EOF
10.0.0.9 ESnew-01
10.0.0.10 ESnew-02
10.0.0.11 ESnew-03
EOF
# 三台都加 hosts,KRaft 需要主机名解析
7)拷贝jdk、kafka、环境变量、hosts
root@ESnew-01 ~# scp -r /usr/local/{jdk8u502-b07,kafka_2.13-3.9.2} ESnew-02:/usr/local
root@ESnew-01 ~# scp -r /usr/local/{jdk8u502-b07,kafka_2.13-3.9.2} ESnew-03:/usr/local
root@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.properties
broker.id=2
# ESnew-02: 2, ESnew-03: 3
advertised.listeners=PLAINTEXT://10.0.0.10:9092
....
root@ESnew-03 ~# source /etc/profile.d/kafka.sh && vim $KAFKA_HOME/config/server.properties
broker.id=3
# ESnew-02: 2, ESnew-03: 3
advertised.listeners=PLAINTEXT://10.0.0.11:9092
root@ESnew-02 ~# mkdir -p /var/lib/kafka
root@ESnew-03 ~# mkdir -p /var/lib/kafka
9)生成集群 UUID
root@ESnew-01 ~# source /etc/profile.d/kafka.sh
root@ESnew-01 ~# kafka-storage.sh random-uuid
xsUe7WA7SXWRxOaDM92JmQ
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 2026
cluster.id=xsUe7WA7SXWRxOaDM92JmQ # 同一集群的 UUID 统一
version=1
directory.id=qqnY7cCi518hXOBd4EfjcQ # 每个节点的唯一目录 ID
node.id=1 # 节点 ID
# ESnew-02: node.id=2, ESnew-03: node.id=3
Caution

⚠️ 三台必须用同一个 UUID 格式化,UUID 是集群的”身份证”

Terminal window
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 集群#

Terminal window
1)创建 topic
root@ESnew-01 ~# kafka-topics.sh --bootstrap-server 10.0.0.9:9092 --topic kpyun-linux --create
Created topic kpyun-linux.
root@ESnew-01 ~# kafka-topics.sh --bootstrap-server 10.0.0.9:9092 --topic kpyun-linux --describe
Topic: 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-beginning
hello-kraft-cluster
✅️ KRaft 集群跨节点消费正常,无需 ZK!
Important

📌 KRaft vs ZK 模式体验对比

对比项ZK 模式KRaft 模式
启动顺序先启 ZK → 再启 Kafka直接启 Kafka
初始化启动即用(ZK 已准备好)先写入集群元数据,再启动
端口只开 9092开 9092 + 9093
配置文件zk.connect 指向 ZKcontroller.quorum.voters

互联网企业架构使用 ELFK 对接 Kafka#

互联网企业架构使用EFLK对接KAFKA
互联网企业架构使用EFLK对接KAFKA

Note

📌 真实企业案例

场景:Ucloud 7 台 Kafka 集群性能差 → 迁移到酒仙桥线下机房 5 台

踩坑记录:

问题现象解决方案
数据丢失酒仙桥专线 2G 打满升级 10G 带宽,HDFS 迁移限速
数据延迟单个 Logstash 高峰期消费慢新增 Logstash 到同消费者组(rebalance 分担)
ES 不显示数据映射(mapping)冲突提前定义索引 mapping

课后练习#

Terminal window
在 ES7 集群,新增一个 elk94 节点,并查看 ES 集群节点列表

文章分享

如果这篇文章对你有帮助,欢迎分享给更多人!

Profile Image of the Author
久棹
不是先学好了再干,而是先干起来再学习,干中学!
分类
站点统计
文章
115
分类
15
标签
272
总字数
306,562
运行时长
0 天
最后活动
0 天前
文章目录