Yuque Archive
寻常的路

Kafka流处理平台(简介)

Kafka结构

目录

  • Kafka概念解析
  • Kafka结构设计
  • Kafka场景与应用
  • Kafka高级特性

Kafka基本概念

  • Producer:消息和数据的生产者,向Kafka的一个topic发布消息的进程/代码/服务
  • Consumer:消息和数据的消费者,订阅数据(Topic)并且处理其发布的消息的进程/代码/服务
  • Consumer Group:逻辑概念,对于同一个topic,会广播给不同的group,一个group中,只有一个consumer可以消费该消息
  • Broker:物理概念,Kafka集群中的每个Kafka节点
  • Topic:逻辑概念,Kafka消息的类别,对数据进行区分、隔离
  • Partition:物理概念,Kafka下数据存储的基本单元。一个Topic数据,会被分散储存到多个Partition,每一个Partition是有序的
  • Replication:同一个Partition可能会有多个Replica,多个Replica之间数据是一样的
  • Replication Leader:一个Partition的多个Replica上,需要一个Leader负责该Partition上与Producer和Consumer交互
  • ReplicaManager:负责管理当前broker所有分区和副本的信息,处理KafkaController发起的一些请求,副本状态的切换、添加/读取消息等

更多Kafka的基本概念

Partition

  • Partition:每一个Topic被切分为多个Partitions
  • 消费者数目少于或等于Partition的数目
  • Broker Group中的每一个broker保存Topic的一个或多个Partitionx
  • Consumer Group中的仅有一个Consumer读取Topic的一个或多个Partitions,并且是惟一的Consumer

Replication

  • 当集群中有Broker挂掉的情况,系统可以主动地使Replicas提供服务
  • 系统默认设置每一个Topic的replication系数为1,可以在创建Topic时单独设置
  • Replication的基本单位是Topic的Partition
  • 所有的读和写都从Leader进,Followers只是做备份
  • Follower必须能够及时复制Leader的数据
  • 增加容错性与可扩展性

Kafka基本结构

  • Producer API
  • Consumer API
  • Strems API
  • Connectors API

Kafka的消息结构

Kafka特点

  • 分布式

  • 多分区

  • 多副本

  • 多订阅者

  • 基于Zookeeper调度

  • 高性能

  • 高吞吐量

  • 低延迟

  • 高并发

  • 时间复杂度为O(1)

  • 持久性与扩展性

  • 数据可持久化

  • 容错性

  • 支持在线水平扩展

  • 消息自动平衡

Kafka应用场景

  • 消息队列
  • 行为跟踪
  • 元信息监控
  • 日志收集
  • 流处理
  • 事件源
  • 持久性日志(commit log)

Kafka高级特性——消息事务

为什么要支持事务

  • 满足“读取-处理-写入”模式
  • 流处理需求的不断增强
  • 不准确的数据处理的容忍度

数据传输的事务定义

  • 最多一次:消息不会被重复发送,最多被传输一次,但也有可能一次不传输
  • 最少一次:消息不会被漏发送,最少被传输一次,但也有可能被重复传输
  • 精确一次:不会被……

事务保证

  • 内部重试问题:Procedure幂等处理
  • 多分区原子写入

事务保证——避免僵尸实例

  • 每个事务Producer分配一个transcational.id,在进程重新启动时能够识别相同的Producer实例
  • Kafka增加了一个与transactional.id相关的epoch,存储每个transactional.id内部元数据
  • 一旦epoch被触发,任何具有相同的transactional.id和…………

Kafka高级特性——零拷贝

简介

  • 网络传输持久性日志块
  • Java Nio channel.transforTo()方法
  • Linux sendfile系统调用

文件传输到网络的公共数据路径

  • 操作系统将数据从磁盘读入到内核空间的页缓存
  • 应用程序将数据从内核空间读入到用户空间缓存中
  • 应用程序将数据写回到内核空间到scoket缓存中
  • 操作系统将数据从socket缓冲区复制到网卡缓冲区,以便将数据经网络发出

零拷贝过程

  • 操作系统将数据从磁盘读入到内核空间的页缓存
  • 将数据的位置和长度的信息的描述符增加至内核空间(scoket缓冲区)
  • 操作系统将数据从内核拷贝到网卡缓冲区,一边讲数据经网络发出

文件传输到网络的公共数据路径演变