至于为什么说它牛叉呢,当然是因为大厂用的比较多嘛!牛逼哄哄的,像什么个性化推荐,微服务,日志处理,应用监控等,都有它的身影!

什么是Kafka
kafka最早诞生于Linkedin的数据管道问题,当时他们采用的是ActiveMQ,但是由于ActiveMQ的各种缺陷导致服务频繁出现故障,所以,当时linkedin的首席架构师lay kreps便组织团队研发了kafka,并随后捐赠给了Apache基金会成为开源项目。
哎,大牛就是大牛啊,咱什么时候才能开发出这么牛叉的软件呢?
kafka简单来说就是一个高吞吐量、分布式的发布-订阅消息系统。其可以很好的处理活跃的流数据,使得数据在各个子系统中高性能、低延迟的不停流转。从而被用在各种实时流数据应用或是对流数据进行响应处理。比如和各种开源分布式系统(如Flume、Apache Storm、Spark、Flink等)集成,提供大数据服务。
Kafka基本结构
在消息队列基础知识那篇文章中,我们介绍过,作为一个消息系统,其基本结构中至少要有产生消息的组件以及消费消息的组件。所以Kafka基本结构如下:

在这里插入图片描述
生产者负责生产消息,将消息写入Kafka集群;消费者从Kafka集群拉取消息;
至于生产者如何将生产的消息写入Kafka,消费者如何从Kafka集群消费消息,Kafka如何存储消息,Kafka集群如何管理调度,如何进行消息负载均衡,以及各组件间如何进行通信等诸多问题,我们会在后面娓娓道来。
Kafka基本概念
1)Topic:主题
Kafka将一组消息抽象为一个主题(Topic),其实就是对消息的一个分类。生产者将消息推送到指定的分类,对应的消费者就会进行消费。
2)Message:消息
消息是Kafka通信的基本单位,由一个固定长度的消息头和一个可变长度的消息体构成。
3)Partition:分区
Kafka将一组消息称为一个主题,而每个主题又被分成一个或多个分区(Partition)。每个分区由一系列有序、不可变的消息组成,是一个有序队列。
从物理角度上,每个分区在机器上对应为一个文件夹,分区的命名规则为主题名称后接“—”连接符,之后再接分区编号,分区编号从0开始,编号最大值为分区的总数减1。每个分区又有一至多个副本(Replica),分区的副本分布在集群的不同代理上,以提高可用性。
从存储角度上分析,分区的每个副本在逻辑上抽象为一个日志(Log)对象,即分区的副本与日志对象是一一对应的。每个主题对应的分区数可以在Kafka启动时所加载的配置文件中配置,也可以在创建主题时指定。当然,客户端还可以在主题创建后修改主题的分区数。
这样的分区使得Kafka在并发处理上变得更加easy!easy!easy!
不过,Kafka只能保证一个分区之内消息的有序性,并不能保证跨分区消息的有序性。每条消息被追加到相应的分区中,是顺序写磁盘,因此效率非常高,这是Kafka高吞吐率的一个重要保证。同时与传统消息系统不同的是,Kafka并不会立即删除已被消费的消息,由于磁盘的限制消息也不会一直被存储(事实上这也是没有必要的),因此Kafka提供两种删除老数据的策略,一是基于消息已存储的时间长度,二是基于分区的大小。
4)Leader副本和Follower副本
由于Kafka副本的存在,就需要保证一个分区的多个副本之间数据的一致性,Kafka会选择该分区的一个副本作为Leader副本,而该分区其他副本即为Follower副本,只有Leader副本才负责处理客户端读/写请求,Follower副本从Leader副本同步数据。如果没有Leader副本,那就需要所有的副本都同时负责读/写请求处理,同时还得保证这些副本之间数据的一致性,假设有n个副本则需要有n×n条通路来同步数据,这样数据的一致性和有序性就很难保证。
引入Leader副本后客户端只需与Leader副本进行交互,这样数据一致性及顺序性就有了保证。Follower副本从Leader副本同步消息,对于n个副本只需n-1条通路即可,这样就使得系统更加简单而高效。副本Follower与Leader的角色并不是固定不变的,如果Leader失效,通过相应的选举算法将从其他Follower副本中选出新的Leader副本。
5)Producer:生产者
生产者(Producer)负责将消息发送给代理,也就是向Kafka代理发送消息的客户端;
6)Comsumer:消费者
消费者(Comsumer)以拉取(pull)方式拉取数据,它是消费的客户端。在Kafka中每一个消费者都属于一个特定消费组(ConsumerGroup),我们可以为每个消费者指定一个消费组,以groupId代表消费组名称,通过group.id配置设置。如果不指定消费组,则该消费者属于默认消费组test-consumer-group。同时,每个消费者也有一个全局唯一的id,通过配置项client.id指定,如果客户端没有指定消费者的id, Kafka会自动为该消费者生成一个全局唯一的id,格式为
{hostName}-
{UUID前8位字符}。同一个主题的一条消息只能被同一个消费组下某一个消费者消费,但不同消费组的消费者可同时消费该消息。消费组是Kafka用来实现对一个主题消息进行广播和单播的手段,实现消息广播只需指定各消费者均属于不同的消费组,消息单播则只需让各消费者属于同一个消费组。
7)Broker:代理
Kafka 集群包含一个或多个服务器,服务器节点称为broker。
broker存储topic的数据。如果某topic有N个partition,集群有N个broker,那么每个broker存储该topic的一个partition。
如果某topic有N个partition,集群有(N+M)个broker,那么其中有N个broker存储该topic的一个partition,剩下的M个broker不存储该topic的partition数据。
如果某topic有N个partition,集群中broker数目少于N个,那么一个broker存储该topic的一个或多个partition。在实际生产环境中,尽量避免这种情况的发生,这种情况容易导致Kafka集群数据不均衡。
应用场景
消息系统是当前处理大数据一个非常重要的组件,用来解决应用解耦、异步通信、流量控制等问题,从而构建一个高效、灵活、消息同步和异步传输处理、存储转发、可伸缩和最终一致性的稳定系统。当前比较流行的消息中间件有Kafka、RocketMQ、RabbitMQ、ZeroMQ、ActiveMQ、MetaMQ、Redis等,这些消息中间件在性能及功能上各有所长。如何选择一个消息中间件取决于我们的业务场景、系统运行环境、开发及运维人员对消息中件间掌握的情况等。大叔认为在下面这些场景中,Kafka是一个不错的选择。
消息系统
Kafka作为一款优秀的消息系统,具有高吞吐量、内置的分区、备份冗余分布式等特点,为大规模消息处理提供了一种很好的解决方案;
应用监控
利用Kafka采集应用程序和服务器健康相关的指标,如CPU占用率、IO、内存、连接数、TPS、QPS等,然后将指标信息进行处理,从而构建一个具有监控仪表盘、曲线图等可视化监控系统。例如,很多公司采用Kafka与ELK(ElasticSearch、Logstash和Kibana)整合构建应用服务监控系统;
用户行为跟踪
为了更好地了解用户行为、操作习惯,改善用户体验,进而对产品升级改进,将用户操作轨迹、内容等信息发送到Kafka集群上,通过Hadoop、Spark或Strom等进行数据分析处理,生成相应的统计报告,为推荐系统推荐对象建模提供数据源,进而为每个用户进行个性化推荐;
流处理
需要将已收集的流数据提供给其他流式计算框架进行处理,用Kafka收集流数据是一个不错的选择,而且当前版本的Kafka提供了Kafka Streams支持对流数据的处理;
持久性日志
Kafka可以为外部系统提供一种持久性日志的分布式系统。日志可以在多个节点间进行备份,Kafka为故障节点数据恢复提供了一种重新同步的机制。同时,Kafka很方便与HDFS和Flume进行整合,这样就方便将Kafka采集的数据持久化到其他外部系统;
总结
这期文章我们主要介绍了一下业界著名的消息队列Kafka,以及Kafka的基本概念、基本组成、应用场景等,其在大厂应用范围极广,也是进入大厂的必备板砖之一,所以大家一定要掌握。
当然,这里做的只是一个简单的入门,更深入的还是需要你去不断的研究和学习。冰冻三尺非一日之寒,技术的成长也非一日之功,需要不断的积累才能强大。
OK,这期文章就写到这里了!
如果你觉得大叔写得还不错,求关注、求点赞、求分享,毕竟大叔熬肝码字也是很幸苦的嘛,当然,如果你是个妹子,也欢迎直接小窗找我哦!嘿嘿嘿……
我是诗远君,一个流浪在互联网江湖的大叔!
我们下期见!




