kafka 简介与部署

环境信息

使用的 hadoop 完全分布式集群, 之前已经安装好 zookeeper

1
2
3
192.168.2.241 hadoop01
192.168.2.242 hadoop02
192.168.2.243 hadoop03

kafka 安装

官网 https://kafka.apache.org/downloads

所有机器操作

1
2
3
4
5
6
7
8
9
wget --no-check-certificate https://archive.apache.org/dist/kafka/3.2.0/kafka_2.13-3.2.0.tgz

useradd kafka
mkdir -p /opt/bigdata/kafka
tar -zxf kafka_2.13-3.2.0.tgz -C /opt/bigdata/kafka
cd /opt/bigdata/kafka/
ln -s kafka_2.13-3.2.0 current

chown -R kafka:kafka /opt/bigdata/kafka/

配置 kafka 集群 (with zookeeper)
/opt/bigdata/kafka/current/config/server.properties

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
# 每个节点不同 eg 1, 2, 3
broker.id=1
listeners=PLAINTEXT://hadoop01:9092 # hadoop02:9092 ,hadoop03:9092 每个节点不同
# 分区数据目录。别指到安装目录里面:
# 一是 kafka-run-class.sh 默认把 log4j 的 LOG_DIR 也设成 $base_dir/logs,
# 分区数据会和 server.log 混在一个目录;
# 二是这条路径穿过 current 软链落在版本目录内,按"换版本只改软链接"的做法升级会把数据留在旧目录。
# 生产建议指向独立数据盘,多盘用逗号分隔。
log.dirs=/data1/kafka,/data2/kafka
num.partitions=6
log.retention.hours=60
log.segment.bytes=1073741824
zookeeper.connect=hadoop01:2181,hadoop02:2181,hadoop03:2181
auto.create.topics.enable=true
delete.topic.enable=true

依次启动

1
2
3
4
$ cd /opt/bigdata/kafka/current
$ nohup bin/kafka-server-start.sh config/server.properties &
$ jps
21840 Kafka

配置 kafka 集群 (without zookeeper) (两选一)
/opt/bigdata/kafka/current/config/kraft/server.properties

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
process.roles=broker,controller
# 每个节点不同 eg 1, 2, 3
node.id=1
controller.quorum.voters=1@hadoop01:19091,2@hadoop02:19091,3@hadoop03:19091
listeners=PLAINTEXT://:9092,CONTROLLER://:19091
inter.broker.listener.name=PLAINTEXT
advertised.listeners=PLAINTEXT://:9092
controller.listener.names=CONTROLLER
listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,SSL:SSL,SASL_PLAINTEXT:SASL_PLAINTEXT,SASL_SSL:SASL_SSL
num.network.threads=3
num.io.threads=8
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600
log.dirs=/opt/bigdata/kafka/current/logs/kraft-combined-logs
num.recovery.threads.per.data.dir=1
offsets.topic.replication.factor=3
# 三节点集群这两项别留单机默认值:单副本时任一 broker 挂掉就丢事务状态、EOS 直接失效
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=2
log.retention.hours=168
log.segment.bytes=1073741824
log.retention.check.interval.ms=3000

启动 kafka

1
2
3
4
5
6
7
8
$ cd /opt/bigdata/kafka/current
$ ./bin/kafka-storage.sh random-uuid # 生成集群 ID
xtzWWN4bTjitpL3kfd9s5g
$ ./bin/kafka-storage.sh format -t xtzWWN4bTjitpL3kfd9s5g -c ./config/kraft/server.properties # 格式化存储目录 所有节点执行

$ nohup ./bin/kafka-server-start.sh ./config/kraft/server.properties & # 启动 kafka 所有节点执行
$ jps
55535 Kafka

验证

创建拥有 3个副本,3 个分区的 topic testtopic

1
2
3
$ ./kafka-topics.sh --bootstrap-server hadoop01:9092,hadoop02:9092,hadoop03:9092 --create -replication-factor 3 --partitions 3 --topic testtopic

Created topic testtopic.

显示 topic

1
2
$ ./kafka-topics.sh --bootstrap-server hadoop01:9092,hadoop02:9092,hadoop03:9092 --list
testtopic

查看 topic testtopic 详细信息

1
2
3
4
5
6
$  ./kafka-topics.sh --bootstrap-server hadoop01:9092,hadoop02:9092,hadoop03:9092 --describe --topic testtopic

Topic: testtopic TopicId: 1o25WxrxTtiswG0nkNf6gw PartitionCount: 3 ReplicationFactor: 3 Configs: segment.bytes=1073741824
Topic: testtopic Partition: 0 Leader: 1 Replicas: 1,3,2 Isr: 1,3,2
Topic: testtopic Partition: 1 Leader: 2 Replicas: 2,1,3 Isr: 2,1,3
Topic: testtopic Partition: 2 Leader: 3 Replicas: 3,2,1 Isr: 3,2,1

生成消息

1
2
3
4
5
6
./kafka-console-producer.sh --broker-list hadoop01:9092,hadoop02:9092,hadoop03:9092 --topic testtopic

### 暂时不输入,等消费信息启动后输入
>hello world
>test kafka
>end kafka

消费信息

1
2
3
4
5
./kafka-console-consumer.sh  --bootstrap-server  hadoop01:9092,hadoop02:9092,hadoop03:9092 --topic testtopic

hello world
test kafka
end

删除信息

1
2
3
4
$ ./kafka-topics.sh --bootstrap-server hadoop01:9092,hadoop02:9092,hadoop03:9092 --delete --topic testtopic
$ ./kafka-topics.sh --bootstrap-server hadoop01:9092,hadoop02:9092,hadoop03:9092 --list

__consumer_offsets

acks 与 ISR:不丢数据的完整配方

上面 producer 配了 acks=all,后面 filebeat 那份配的是 required_acks: 1。同一条采集链路两端不一致,正好可以拿来说清这件事——因为光有 acks=all 并不保证不丢数据

先看 acks 三个取值:

取值 语义 丢数据的条件
0 发出去就算成功,不等任何响应 网络抖一下就丢,吞吐最高
1 leader 写进自己的日志就返回 leader 落盘后、副本还没跟上就宕机 → 丢
all / -1 当前 ISR 里所有副本都确认才返回 见下

关键在 all 那一行的”当前 ISR”。ISR(In-Sync Replicas)是动态收缩的:副本落后超过 replica.lag.time.max.ms(默认 30s)就被踢出去。于是有这样一条路径:

  1. 三副本,正常时 ISR = {leader, f1, f2};
  2. 两个 follower 都因为 GC、磁盘慢或网络问题掉队,被踢出 ISR;
  3. 此时 ISR 只剩 leader 一个acks=all 就等价于 acks=1——它”所有副本都确认了”,只不过所有副本就是它自己;
  4. leader 这时宕机,那批只写进 leader 的数据就没了。

所以真正的配方是四项一起配,缺一不可:

1
2
3
4
5
6
7
# broker 端
replication.factor=3 # 三副本
min.insync.replicas=2 # ISR 少于 2 个就拒绝写入(报 NotEnoughReplicas)
unclean.leader.election.enable=false # 不允许落后的副本当 leader

# producer 端
acks=all

min.insync.replicas=2 是补上第 3 步那个漏洞的那一块:ISR 缩到 1 时,写入直接失败而不是”看起来成功了”。这是一个明确的取舍——用可用性换持久性,宁可让 producer 报错重试,也不接受静默丢数据。

unclean.leader.election.enable 这一项被违反的后果更严重:允许一个不在 ISR 里的、数据落后的副本当选 leader,等于丢弃已经提交过的数据,而且消费者那边会看到 offset 回退。这个参数在较新版本里默认已经是 false,但升级上来的老集群里常有 true 的历史包袱,值得单独检查一遍。

三副本 + min.insync.replicas=2 的组合还有个好处:允许一台 broker 计划内下线(滚动重启、打补丁)而不影响写入,因为剩下两个仍满足最小同步副本数。

至于上面 filebeat 的 required_acks: 1——日志采集这种场景丢几条通常可以接受,用 1 换吞吐是合理选择。但要意识到这是个明确的降级决定,而不是默认值刚好如此;真正要求不丢的链路两端都得写 -1(filebeat 里 required_acks: -1 对应 acks=all)。

kafka 集群之间同步数据

环境信息

新建 kafka 集群,以及一个 客户端(最好独立运行,也可以放在新建的kafka集群中,当前独立一个节点运行)

所有节点

1
2
3
4
5
6
7
8
9
10
$ cat /etc/hosts
127.0.0.1 localhost localhost.localdomain localhost4 localhost4.localdomain4
::1 localhost localhost.localdomain localhost6 localhost6.localdomain6
192.168.2.171 kafka01
192.168.2.173 kafka03
192.168.2.174 kafka04
192.168.2.241 hadoop01
192.168.2.242 hadoop02
192.168.2.243 hadoop03
192.168.2.86 kafkaclient

其中 kafka01,03,04 这三台为备用集群,hadoop01,02,03 为主要集群, 均已启动

MirrorMaker 的配置

  1. /opt/bigdata/kafka/current/config/consumer.properties
1
2
3
4
5
6
7
8
9
bootstrap.servers=hadoop01:9092,hadoop02:9092,hadoop03:9092
group.id=hadoop
enable.auto.commit=false
request.timeout.ms=180000
heartbeat.interval.ms=1000
session.timeout.ms=120000
max.poll.interval.ms=600000
max.poll.records=120000
auto.offset.reset=earliest
  1. /opt/bigdata/kafka/current/config/producer.properties
1
2
3
4
5
6
7
bootstrap.servers=kafka01:9092,kafka03:9092,kafka04:9092
acks=all
batch.size=16348
linger.ms=1
max.block.ms=9223372036854775807
compression.type=gzip
request.timeout.ms=90000

测试

  1. 启动 MirrorMaker 服务

kafkaclient

1
2
cd /opt/bigdata/kafka/current
nohup bin/kafka-mirror-maker.sh --consumer.config ./config/consumer.properties --num.streams 16 --producer.config ./config/producer.properties --whitelist="mglogs.*" &

-whitelist:设置要同步的 Topic。收的是 Java 正则不是 glob——写 mglogs* 意思是”mglog 后面跟 0 个或多个 s“,会把 mglog 一起同步进来,而 mglogs-2024 反倒不匹配。要前缀匹配得写 mglogs.*,精确匹配就直接写 mglogs

  1. 验证
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
[kafka@host86 bin]$ ./kafka-consumer-groups.sh  --bootstrap-server hadoop01:9092,hadoop02:9092,hadoop03:9092 --describe --group hadoop

GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID
hadoop mglogs 5 - 0 - hadoop-13-8744d13f-d252-459a-b1fb-223f5007d41a /192.168.2.86 hadoop-13
hadoop mglogs 1 1 1 0 hadoop-1-b2ff1495-7c4c-4a49-b4a8-6a36b2873fa2 /192.168.2.86 hadoop-1
hadoop mglogs 3 - 0 - hadoop-11-76351435-b37f-4596-81ed-e8ce1a730715 /192.168.2.86 hadoop-11
hadoop mglogs 0 2 2 0 hadoop-0-c70aa0b9-a2f1-417c-a4c6-dfbd621752f0 /192.168.2.86 hadoop-0
hadoop mglogs 2 - 0 - hadoop-10-7128d38d-1d31-45d9-9845-acc4831d2384 /192.168.2.86 hadoop-10
hadoop mglogs 4 - 0 - hadoop-12-11bac4a8-be82-4c1b-be3a-8e0231008c9c /192.168.2.86 hadoop-12
[kafka@host86 bin]$ ./kafka-console-consumer.sh --bootstrap-server hadoop01:9092,hadoop02:9092,hadoop03:9092 --topic mglogs --from-beginning
May 22 21:51:49 Installed: 2:nmap-ncat-6.40-19.el7.x86_64
May 22 21:51:49 Installed: 14:libpcap-1.5.3-13.el7_9.x86_64
hello world
^CProcessed a total of 3 messages
[kafka@host86 bin]$ ./kafka-console-consumer.sh --bootstrap-server kafka01:9092,kafka03:9092,kafka04:9092 --topic mglogs --from-beginning
hello world
May 22 21:51:49 Installed: 14:libpcap-1.5.3-13.el7_9.x86_64
May 22 21:51:49 Installed: 2:nmap-ncat-6.40-19.el7.x86_64
^[c^CProcessed a total of 3 messages
[kafka@host86 bin]$ ./kafka-topics.sh --bootstrap-server hadoop01:9092,hadoop02:9092,hadoop03:9092 --list
__consumer_offsets
mglogs
my_test
test_topic

[kafka@host86 bin]$ ./kafka-topics.sh --bootstrap-server kafka01:9092,kafka03:9092,kafka04:9092 --list
__consumer_offsets
mglogs
[kafka@host86 bin]$ # 只同步了 mglogs

备注
数据收集采用的 filebeat 输出到 kafka, 具体可参考 ELK 组件部署 中的 filebeat 一节

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
filebeat.inputs:
- type: log
enabled: true
paths:
- /var/log/*.log
fields:
log_topic: mglogs
filebeat.config.modules:
path: ${path.config}/modules.d/*.yml
reload.enabled: false
name: "appserver1"
output.kafka:
enabled: true
hosts: ["hadoop01:9092", "hadoop02:9092", "hadoop03:9092"]
version: "0.10"
topic: '%{[fields][log_topic]}'
codec.format.string: '%{[message]}'
partition.round_robin:
reachable_only: true
worker: 2
required_acks: 1
compression: gzip
max_message_bytes: 10000000
processors:
- drop_fields:
fields: ["input", "host", "agent.type", "agent.ephemeral_id", "agent.id", "agent.version", "ecs"]
logging.level: info