前言
最近一直在写文档,但是也穿插着其他各项工作。比如就有一个客户,整天问TiCDC和Kafka的问题。由于刚刚接触TiCDC不久,对Kafka也是浅尝辄止的弄过。因此经常被客户问懵圈。
用kafka这类软件对DBA发问,大风大浪的我见过世面的,稍微学习一下,立马支棱起来。
了解Kafka
首先我们先科普一下Kafka。kafka的架构如图所示:

其实整个架构有点类似我们日常生活中的什么场景呢?小区门口的菜鸟驿站。
想想那些搞快递的小哥们,是不是只要把快递丢到菜鸟驿站就行了,然后你下班了或者有空,去菜鸟驿站直接取快递就行了。假设没有菜鸟驿站,你和快递小哥之间传递快递,会出现各种各样的问题。一种情况是他问你现在有时间取快递吗?你可能正在上班或者做着重要的事情,当然没时间立马取。而快递小哥也不可能一直在等你。所以菜鸟驿站应运而生。果然程序设计还是来源于生活。
自从有了菜鸟驿站,快递小哥只用专注的把货物送到菜鸟驿站就行了,而咱们也可以下班到点了过去取快递。这样就可以轻松实现业务解耦,异步处理,流量控制等。
回过头来再看上面的架构图,是不是清晰了很多。左侧的应用就叫生产者(Producer),负责往kafka集群扔数据。而最右面的是消费者(Consumer),负责获取数据。而中间的就是Kafka集群,一个集群由多个Broker组成,Broker负责接收和处理客户端发送过来的请求,以及对消息进行持久化。Topic,消息的分类,同一类消息属于同一个Topic。而每个topic又可以有多个Partition(分区)。如上图中的T1P1,T1P2,T1P3。一个topic下的多个分区可以并发接收消息,同样的也能供消费者并发拉取消息。
Raft方式安装Kafka
接下来我们来安装kafka,通过安装边学边了解,我们使用Raft模式来安装,原因就是不想再去弄个ZK。下载3.0的介质然后解压。
wget -c http://dlcdn.apache.org/kafka/3.0.0/kafka_2.13-3.0.0.tgz
tar xvf kafka_2.13-3.0.0.tgz
cd kafka_2.13-3.0.0/
Kafka集群配置
这里我们要在单机环境模拟3个节点的集群,先复制三份配置文件。
cd config/kraft
cp server.properties server1.properties
cp server.properties server2.properties
cp server.properties server3.properties
在server1.properties
中,修改以下属性。然后保持其他属性不变。
node.id=1
process.roles=broker,controller
inter.broker.listener.name=PLAINTEXT
controller.listener.names=CONTROLLER
listeners=PLAINTEXT://:9092,CONTROLLER://:19092
Advertising.listeners = PLAINTEXT://localhost:9092
log.dirs=/home/kafka/server1/kraft-combined-logs
listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,SSL:SSL,SASL_PLAINTEXT:SASL_PLAINTEXT,SASL_SSL:SASL_SSL
controller.quorum.voters=1@localhost:19092,2@localhost:19093,3@localhost:19094
先解释一下这些属性的作用:
node.id:这将作为集群中的节点ID。
process.roles :指明节点既可以充当broker,也可以充当controller。
前面说过broker,那controller是什么意思呢 ?
简单来说,就是在分布式系统中,通常需要有一个协调者,在kafka中这个协调者就是控制器(Controller),它本身其实也是一个Broker,只不过需要负责一些额外的工作(追踪集群中的其他Broker,并在合适的时候处理新加入的和失败的Broker节点、Rebalance分区、分配新的leader分区等)。
inter.broker.listener.name :该参数是用来设置broker之间进行通信时采用的listener名称,一般来说默认为PLAINTEXT。
controller.listener.names : 这里控制器监听器名称设置为CONTROLLER
listeners : 指定broker使用 9092 端口,而kraft controller将使用 19092 端口
advertised.listeners : 和 listeners 相比多了个 advertised。表示宣称的、公布的,就是说这组监听器是用于对外发布的。
log.dirs:指定kafka的数据日志目录。
listener.security.protocol.map :将监听名称映射到安全协议,默认情况下它们是相同的
controller.quorum.voters :这个配置标识出有哪些节点是 Quorum的投票者节点,所有想成为控制器的节点都需要包含在这个配置里面。
对于server2.properties和server3.properties修改以下属性。保持其他属性不变。
node.id=2
process.roles=broker,controller
inter.broker.listener.name=PLAINTEXT
controller.listener.names=CONTROLLER
listeners=PLAINTEXT://:9093,CONTROLLER://:19093
Advertising.listeners = PLAINTEXT://localhost:9093
log.dirs=/home/kafka/server2/kraft-combined-logs
listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,SSL:SSL,SASL_PLAINTEXT:SASL_PLAINTEXT,SASL_SSL:SASL_SSL
controller.quorum.voters=1@localhost:19092,2@localhost:19093,3@localhost:19094
node.id=3
process.roles=broker,controller
inter.broker.listener.name=PLAINTEXT
controller.listener.names=CONTROLLER
listeners=PLAINTEXT://:9094,CONTROLLER://:19094
Advertising.listeners = PLAINTEXT://localhost:9094
log.dirs=/home/kafka/server3/kraft-combined-logs
listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,SSL:SSL,SASL_PLAINTEXT:SASL_PLAINTEXT,SASL_SSL:SASL_SSL
controller.quorum.voters=1@localhost:19092,2@localhost:19093,3@localhost:19094
上述配置完成之后,我们首先要创建集群🆔。
./bin/kafka-storage.sh random-uuid
lxayDXaISrO1b_XiQ-i9zw
接下来我们需要格式化所有的存储目录。也就是我们放在log.dirs属性中的目录位置。
./bin/kafka-storage.sh format -t <uuid> -c <server_config_location>
./bin/kafka-storage.sh format -t lxayDXaISrO1b_XiQ-i9zw -c ./config/kraft/server1.properties
./bin/kafka-storage.sh format -t lxayDXaISrO1b_XiQ-i9zw -c ./config/kraft/server2.properties
./bin/kafka-storage.sh format -t lxayDXaISrO1b_XiQ-i9zw -c ./config/kraft/server3.properties
然后就可以启动 kafka 服务器。
export KAFKA_HEAP_OPTS="-Xmx200M –Xms100M"
./bin/kafka-server-start.sh -daemon ./config/kraft/server1.properties
./bin/kafka-server-start.sh -daemon ./config/kraft/server2.properties
./bin/kafka-server-start.sh -daemon ./config/kraft/server3.properties
创建主题
安装好kafka之后,我们就可以创建主题了。
./bin/kafka-topics.sh --create --topic kafka-test --partitions 3 --replication-factor 3 --bootstrap-server localhost:9092
由于我们有 3 个节点,因此我们创建一个具有 3 个分区和 3 个副本的主题。
列出集群中存在的主题
./bin/kafka-topics.sh --bootstrap-server localhost:9093 --list
kafka-test
查看更加详细的主题信息
./bin/kafka-topics.sh --bootstrap-server localhost:9093 --describe --topic kafka-test
Topic: kafka-test TopicId: pITLMOK9RyioQLpSr11XmQ PartitionCount: 3 ReplicationFactor: 3 Configs: segment.bytes=1073741824
Topic: kafka-test Partition: 0 Leader: 3 Replicas: 3,1,2 Isr: 3,1,2
Topic: kafka-test Partition: 1 Leader: 1 Replicas: 1,2,3 Isr: 1,2,3
Topic: kafka-test Partition: 2 Leader: 2 Replicas: 2,3,1 Isr: 2,3,1
生产和消费
接下来我们来模拟程序的生产和消费。我们先打开一个客户端,充当生产者。
./bin/kafka-console-producer.sh --broker-list localhost:9092 --topic kafka-test
然后再打开另外一个客户端,充当消费者。
./bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic kafka-test
在生产者产生消息。

然后消费者你将会看到这些消息。

关闭集群
最后玩完了,可以关闭集群,由于一台机器部署了三个kafka,使用stop竟然一下可以停掉三个,有点神奇。
./bin/kafka-server-stop.sh stop
后记
今天是kafka的一道前菜,了解基本概念和安装。如果不懂为什么需要有kafka这个东西,那么还是从具体的业务上来了解比较好,可以找你身边的业务人员问问。比如现在一个业务场景就是要对一个产品进行报价,我们通过一个报价系统生成报价,然后我们把报价发送到Kafka上就行了。剩下的下游应用如果想获取报价,就去订阅信息就行了。
最后是一个思考人生的问题,DBA究竟要不要学会kafka ?
这其实和你接触的场景有关,如果你接触的事情很多都是CDC复制场景,我建议还是支棱起来,始终是绕不开kafka的。
[参考文档]: https://juejin.cn/post/6844903860494942222




