环境信息
使用的 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)就被踢出去。于是有这样一条路径:
- 三副本,正常时 ISR = {leader, f1, f2};
- 两个 follower 都因为 GC、磁盘慢或网络问题掉队,被踢出 ISR;
- 此时 ISR 只剩 leader 一个,
acks=all 就等价于 acks=1——它”所有副本都确认了”,只不过所有副本就是它自己;
- leader 这时宕机,那批只写进 leader 的数据就没了。
所以真正的配方是四项一起配,缺一不可:
1 2 3 4 5 6 7
| replication.factor=3 # 三副本 min.insync.replicas=2 # ISR 少于 2 个就拒绝写入(报 NotEnoughReplicas) unclean.leader.election.enable=false # 不允许落后的副本当 leader
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 的配置
- /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
|
- /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
|
测试
- 启动 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 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
|