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
相关产品推荐
相关产品推荐

