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

Kafka Connect构建SourceRecord时如何处理自定义对象List字段

问题根因

抛出Invalid Java object for schema type STRUCT: class com.dto.Currencies异常的核心原因是:Kafka Connect的Struct体系不识别自定义Java POJO类,无论单个STRUCT类型字段,还是ARRAY类型里嵌套的STRUCT元素,都必须传入Connect API定义的Struct实例,不能直接传入自行编写的Currencies DTO对象。

解决方案

1. 修正Schema声明

原有Schema的结构逻辑没有大问题,唯一需要调整的是数组内部嵌套的STRUCT必须显式指定名称,否则嵌套结构在序列化校验时会报错,修正后的Schema代码如下:

public static final Integer FIRST_VERSION = 1;
public static final String CURRENCIES_FIELD_NAME = "currencies";
public static final String CURRENCY_ENTITY_NAME = "com.dto.Currencies";

// 先单独定义数组内单个Currencies对象的STRUCT Schema
public static final Schema CURRENCY_STRUCT_SCHEMA = SchemaBuilder.struct()
        .name(CURRENCY_ENTITY_NAME)
        .version(FIRST_VERSION)
        .field("code", Schema.OPTIONAL_STRING_SCHEMA)
        .field("title", Schema.OPTIONAL_STRING_SCHEMA)
        .field("slug", Schema.OPTIONAL_STRING_SCHEMA)
        .field("url", Schema.OPTIONAL_STRING_SCHEMA)
        .optional()
        .build();

// 再定义currencies字段的数组Schema
public static final Schema CURRENCIES_ARRAY_SCHEMA = SchemaBuilder.array(CURRENCY_STRUCT_SCHEMA)
        .optional()
        .name(CURRENCIES_FIELD_NAME)
        .version(FIRST_VERSION)
        .build();

// 外层新闻Schema
public static final Schema NEWS_SCHEMA = SchemaBuilder.struct()
        .name("News")
        .version(FIRST_VERSION)
        .field(CURRENCIES_FIELD_NAME, CURRENCIES_ARRAY_SCHEMA)
        // 其余普通字段按原有逻辑定义即可
        .build();

2. 修正Struct赋值逻辑

赋值时不能直接传入List<Currencies> POJO列表,需要遍历列表中的每一个Currencies对象,将其转换为对应Schema的Struct实例,再把转换后的Struct列表传入外层Struct,修正后的赋值代码如下:

public Struct buildRecordValue(CryptoNews cryptoNews){
    Struct valueStruct = new Struct(NEWS_SCHEMA);

    List<Currencies> currencies = cryptoNews.getCurrencies();
    if (currencies != null) {
        List<Struct> currencyStructList = new ArrayList<>();
        // 逐个将POJO转换为Connect Struct对象
        for (Currencies currency : currencies) {
            Struct currencyStruct = new Struct(CURRENCY_STRUCT_SCHEMA);
            currencyStruct.put("code", currency.getCode());
            currencyStruct.put("title", currency.getTitle());
            currencyStruct.put("slug", currency.getSlug());
            currencyStruct.put("url", currency.getUrl());
            currencyStructList.add(currencyStruct);
        }
        // 传入转换好的Struct列表
        valueStruct.put(CURRENCIES_FIELD_NAME, currencyStructList);
    }

    // 其余普通字段按原有逻辑赋值即可
    return valueStruct;
}

注意:当前使用的0.10.2.0版本Kafka Connect对嵌套Schema名校验较严格,不要省略嵌套STRUCT的name配置,否则会出现Schema不匹配的序列化错误。现有worker.properties里的JsonConverter配置不需要额外修改,上述逻辑可以直接适配。

内容的提问来源于stack exchange,提问作者leonaugust

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.02 06:01:43