GBase开启kafka transaction topic后,无法正确解析到kafka传来的数据
【GBase版本】: 9.5.2
【操作系统】: centos 7
【问题描述】*: 实验目的是打通Oracle到Kafka到GBase 8a数据同步的整条链路,目前kafka端能正确接收到数据,其格式为json;GBase端创建kafka transaction consumer,且其状态为开启状态,但GBase中无法同步到kafka中的数据信息。使用语句select exception from information_schema.kafka_consumer_status;
得到如下报错信息:

字面意思应该是这两个参数的配置问题,请问在哪里设置这两个参数的值呢?
评论


如果您确实是按照上述规则创建的,修改参数方法如下:
set global gcluster_kafka_max_message_size = 1000000000

针对这个报错分析,造成该问题的两个参数含义如下:
receive.message.max.bytes: kafka允许的最大消息集合批次,生产者往kafka发送消息的一批消息最大不能超过这个参数。
replica.fetch.max.bytes:限制拉取分区中消息的大小。该值必须大于message.max.bytes,否则会因replica.fetch.max.bytes过小,导致数据不同步,备份不成功的报错问题。
该报错信息体现的就是,kafka配置了大消息的同时,8a集群作为客户端无法创建rdkafka(kafka consumer),理由是8a中配置的最大接收消息大小太小。
在8a中,receive.message.max.bytes体现为参数,gcluster_kafka_max_message_size。
该参数含义为:从kafka topic获得消息的最大长度,单位为字节。
最小值:1000、最大值:1000000000、建议值:100000000
同时与该参数存在关联关系的参数是,gcluster_kafka_fetch_max_size,配置要求为:
gcluster_kafka_max_message_size >= gcluster_kafka_fetch_max_size + 512
用于协议的额外开销。
以上两个参数可使用set语句修改值,也可以在配置文件中修改,适用于session、global范围均可。
oracle中待同步的表test有id和name两列,id是主键;GBase端的建表语句是create table test (id int, name varchar(20))。请问发生此错误的原因是两个表未能对应吗?还是kafka中的json数据存在一些问题呢?



GBase 8a kafka consumer要求json格式的数据需要包含主键,如"primary_key":{"A","B"},或者不带主键的普通文本


可以尝试将GBase端的表也建为带有主键的表。
json消息格式举例如下:
{ "table":"BDTEST.TEST4", "op_type":"I", "op_ts":"2022-01-16 09:26:29.707674", "current_ts":"2022-01-16T17:26:34.556001", "pos":"00000000030000002194", "after":{ "A":4, "B":40, "C":"t4" } }
{ "table":"BDTEST.TEST4", "op_type":"D", "op_ts":"2022-01-16 09:36:44.703860", "current_ts":"2022-01-16T17:36:49.047000", "pos":"00000000030000003188", "primary_keys":{"A"}, "before":{ "A":20 } }
{ "table":"BDTEST.TEST4", "op_type":"U", "op_ts":"2022-01-16 09:32:33.705303", "current_ts":"2022-01-16T17:32:36.839000", "pos":"00000000030000002612", "primary_keys":{"A"}, "before":{ "A":2 } "after":{ "A":20, "B":200, "C":"t20" } }
热门帖子
- 12025-12-01浏览数:182759
- 22023-05-09浏览数:25051
- 42023-09-25浏览数:18521
- 52020-05-11浏览数:17526
目前kafak中有数据,您只需要消费kafka数据就可以,创建kafka consumer对kafka数据消费
语法:create kafka consumer consumer_name transaction topic topic_name brokers 'broker_IP:9092,broker_IP:9092';
创建之后启动kafka consumer
start kafka consumer consumer_name;