未使用Schema Registry的Kafka Connect BigQuery连接器写入报错求解
问题引发原因
核心报错根因为Could not convert to BigQuery schema with a batch of tombstone records,触发逻辑如下:
- 你启用了无Schema模式的BigQuery Sink Connector,连接器启动后首次创建目标表时,需要用第一批收到的消息字段自动推导BigQuery表结构
- 你当前第一批到达连接器的消息全是墓碑记录(即Kafka中value为null的消息,通常用于compact主题的删除标记),没有有效业务字段,连接器无法从全空的value样本中推导出合法的表结构,因此抛出异常
- 你新增的
allowBigQueryRequiredFieldRelaxation、allBQFieldsNullable两个配置仅用于处理BigQuery字段非空约束兼容问题,和本次无有效Schema可推导的报错场景无关,因此无法解决问题
修复方案
- 方案一:提前手动建表
关闭连接器自动建表能力,提前在BigQuery中创建好对应业务结构的目标表,配置新增:
连接器不需要自动推导Schema即可直接写入数据,不会触发该报错autoCreateTables: false - 方案二:过滤墓碑消息
通过Kafka Connect自带的转换组件直接过滤掉value为null的墓碑记录,配置新增:
所有空value的墓碑消息会被直接过滤,不会进入Schema推导和写入流程transforms: dropTombstones transforms.dropTombstones.type: org.apache.kafka.connect.transforms.Filter$Value transforms.dropTombstones.predicate: isTombstone predicates: isTombstone predicates.isTombstone.type: org.apache.kafka.connect.predicates.NullValue - 方案三:调整消息写入顺序
如果你的业务场景需要保留墓碑消息不需要过滤,先向目标Topic写入至少一条有有效value的非空业务消息,等连接器成功完成自动建表后,再写入包含墓碑消息的业务数据即可正常运行
内容的提问来源于stack exchange,提问作者S S
相关产品推荐
相关产品推荐

