You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.03 07:24:03