GBase 8a kafka consumer 功能介绍
8a kafka consumer
Kafka consumer 的主要功能就是同步 Kafka 数据到 GBase 8a MPP Cluster:根据配置,可以指定需要同步的业务; 在同步过程中,提供同步状态查询功能; 实现数据同步的高可用性和事务数据一致性。
什么是kafka
概念:
Kafka是由Apache软件基金会开发的一个开源流处理平台,由Scala和Java编写。Kafka是一种高吞吐量的分布式发布订阅消息系统, Kafka最早设计的目的是作为LinkedIn的活动流和运营数据的处理管道,这些数据主要是用来对用户做画像分析以及服务器性能数据的监控,所以Kafka一开始设计的目标就是作为一个分布式、高吞吐量的消息系统,所以适合运用在大数据传输场景。
特性:
通过O(1)的磁盘数据结构提供消息的持久化,这种结构对于即使数以TB的消息存储也能够保持长时间的稳定性能
高吞吐量 :即使是非常普通的硬件Kafka也可以支持每秒数百万 的消息。
支持通过Kafka服务器和消费机集群来分区消息。
支持Hadoop并行数据加载。
kafka一些概念
Broker : 和AMQP里协议的概念一样, 就是消息中间件所在的服务器
Topic(主题) : 每条发布到Kafka集群的消息都有一个类别,这个类别被称为Topic。(物理上不同Topic的消息分开存储,逻辑上一个Topic的消息虽然保存于一个或多个broker上但用户只需指定消息的Topic即可生产或消费数据而不必关心数据存于何处)
Partition(分区) : Partition是物理上的概念,体现在磁盘上面,每个Topic包含一个或多个Partition.
Producer : 负责发布消息到Kafka broker
Consumer : 消息消费者,向Kafka broker读取消息的客户端。
Consumer Group(消费者群组) : 每个Consumer属于一个特定的Consumer Group(可为每个Consumer指定group name,若不指定group name则属于默认的group)。
offset 偏移量: 是kafka用来确定消息是否被消费过的标识,在kafka内部体现就是一个递增的数字
Broker : 和AMQP里协议的概念一样, 就是消息中间件所在的服务器
事务型Consumer 具体做什么
目标:准确无误地把kafka队列里的消息(数据)同步到8A
消费和解析
•调用kafka client lib,从kafka队列消费消息
•初步扫描消息,判断格式是否存在问题
•解析消息,得到DML操作所需的要素(db、table、primary key、insert values、delete where)
•获取8A端的表结构
•验证primary key的一致性
•验证源表表定义有没有出现变化(因为现在尚未支持DDL同步)
•验证源表、目标表表定义的一致性(旧版要求完全一致,新版允许源表是目标表的子集)
数据同步
•以事务方式执行DML操作
•根据配置参数决定何时提交
•提交时保存当前消费进度
•支持高可用
自动重试
事务型consumer消费json格式说明
json消息的编码,采用UTF8.
一个json消息包含1个或多个{},{},片段,每个片段代表一行数据操作,例如
{"table":"db.table","op_type":"I","POS":"10000000990001203450","primary_keys":["A","B"],"after":{"A":"abcd","B":100,"C":"ddd"}},
如果一个json消息包含多个{},片段,则kafka consumer将其视为一个事务内的若干操作,会保证其原子性。
热门帖子
- 12025-12-01浏览数:182759
- 22023-05-09浏览数:25052
- 42023-09-25浏览数:18521
- 52020-05-11浏览数:17526