GBase 8a
运维管理
文章

GBase 8a kafka consumer 功能介绍

发表于2024-12-30 18:25:45249次浏览1个评论

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将其视为一个事务内的若干操作,会保证其原子性。

评论

登录后才可以发表评论
崔哥发表于 5个月前
兔园标物序,惊时最是梅。衔霜当路发,映雪拟寒开。枝横却月观,花绕凌风台。朝洒长门泣,夕驻临邛杯。应知早飘落,故逐上春来。