Kafka Connect JDBC:自定义SMT还是链式SMT可简化嵌套Avro载荷?
问题原因说明
你之前用Flatten转换没有效果,是因为Flatten SMT的作用是合并嵌套的字段路径(比如将payload.id.int转为payload_id_int),但你当前的结构是每个字段的取值本身是仅含单个类型键的对象,不属于字段路径嵌套的场景,所以Flatten不会对这部分做处理。
以下是按优先级排序的可行解决方案:
方案1:调整JDBC源连接器配置(最高效,最推荐)
你遇到的嵌套结构是JDBC Source Connector针对允许为NULL的字段,用Avro序列化时默认生成的联合类型包装结构,直接在连接器配置中增加以下参数即可从源头生成扁平结构,不需要额外做任何转换:
# 关闭Avro转换器的元数据包装 value.converter.connect.meta.data=false # 禁止JDBC源默认将字段设为允许为空,避免联合类型生成 jdbc.source.default.nulls.allowed=false # 若使用JSON转换器则补充这个配置 value.converter.schemas.enable=false
配置后重启连接器,新写入Kafka Topic的消息payload直接就是你需要的扁平格式。
方案2:自定义轻量SMT转换(适合不能修改源连接器配置的场景)
如果不能调整源连接器的配置,写一个几十行代码的自定义Kafka Connect SMT即可实现全字段自动打平,不需要为每个字段单独配置转换规则,100个字段也能自动处理,性能损耗极低。核心逻辑参考:
// 省略类定义和配置部分,仅保留核心转换逻辑 @Override public R apply(R record) { Struct originalPayload = (Struct) record.value().get("payload"); SchemaBuilder newSchemaBuilder = SchemaBuilder.struct(); // 遍历原字段构建新Schema for (Field field : originalPayload.schema().fields()) { Schema nestedSchema = originalPayload.getStruct(field.name()).schema(); newSchemaBuilder.field(field.name(), nestedSchema.fields().get(0).schema()); } Schema newSchema = newSchemaBuilder.build(); Struct newPayload = new Struct(newSchema); // 填充字段值 for (Field field : originalPayload.schema().fields()) { Struct nestedVal = originalPayload.getStruct(field.name()); newPayload.put(field.name(), nestedVal.get(nestedVal.schema().fields().get(0))); } // 返回新的记录 return record.newRecord( record.topic(), record.kafkaPartition(), record.keySchema(), record.key(), newSchema, newPayload, record.timestamp() ); }
打好JAR包放到Kafka Connect的插件目录,在连接器配置中启用该SMT即可。
方案3:KSQL流转换(适合不想写代码的临时场景)
如果不想开发自定义SMT,也可以用KSQL做一层中转转换,将原始Topic的数据处理后写入新的扁平结构Topic,示例语句:
-- 关联原始Topic创建流 CREATE STREAM customers_raw ( payload STRUCT< id STRUCT<int INT>, full_name STRUCT<string STRING>, birthdate STRUCT<int INT>, fav_animal STRUCT<string STRING>, fav_colour STRUCT<string STRING>, fav_movie STRUCT<string STRING> > ) WITH (KAFKA_TOPIC='mysql-00-customers', VALUE_FORMAT='AVRO'); -- 生成扁平结构的新流 CREATE STREAM customers_flat WITH ( KAFKA_TOPIC='mysql-00-customers-flat', VALUE_FORMAT='AVRO', PARTITIONS=6 ) AS SELECT payload->id->int AS id, payload->full_name->string AS full_name, payload->birthdate->int AS birthdate, payload->fav_animal->string AS fav_animal, payload->fav_colour->string AS fav_colour, payload->fav_movie->string AS fav_movie FROM customers_raw EMIT CHANGES;
该方案的缺点是字段较多时写查询语句比较繁琐,还会多占用一份Topic存储资源。
内容的提问来源于stack exchange,提问作者Cr4zyTun4
相关产品推荐
相关产品推荐

