能否使用Kafka Transforms InsertField为kafka-connect-solr消息嵌套add字段避免Solr覆盖
结论
仅靠内置的InsertField Transform 无法实现你需要的消息结构转换,更推荐自定义Kafka Connect单消息转换(SMT)实现需求。
原因说明
内置InsertField的核心能力仅为新增固定值/上下文变量/复制现有字段原始值的新字段,不支持将已有字段的标量值包装为{"add": [原值]}的嵌套结构,也无法直接修改原有字段的结构类型,完全匹配不了你的转换需求。
落地可选方案
方案1:自定义SMT(最优推荐)
代码实现逻辑非常简单,性能损耗可忽略,完全适配大数据量场景:
- 仅针对Topic2的消息做处理,保留
id字段原样 - 遍历所有非
id的业务字段,将字段的标量值统一包装为{"add": [原值]}的嵌套结构 - 打包为jar包放到Kafka Connect插件目录后,直接在connector配置中引用即可
核心处理逻辑参考:
@Override public R apply(R record) { Object originalVal = operatingValue(record); if (originalVal == null || !record.topic().equals("你的Topic2名称")) { return record; } Struct originalStruct = (Struct) originalVal; // 构建新的消息结构,动态处理所有非id字段 SchemaBuilder newSchemaBuilder = SchemaBuilder.struct().field("id", originalStruct.schema().field("id").schema()); Map<String, Object> fieldValues = new HashMap<>(); fieldValues.put("id", originalStruct.get("id")); for (Field field : originalStruct.schema().fields()) { if ("id".equals(field.name())) continue; // 构建add嵌套结构 Schema addSchema = SchemaBuilder.struct() .field("add", SchemaBuilder.array(field.schema()).optional().build()) .optional().build(); newSchemaBuilder.field(field.name(), addSchema); Struct addStruct = new Struct(addSchema); addStruct.put("add", Collections.singletonList(originalStruct.get(field.name()))); fieldValues.put(field.name(), addStruct); } Struct newStruct = new Struct(newSchemaBuilder.build()); fieldValues.forEach(newStruct::put); return newRecord(record, newSchemaBuilder.build(), newStruct); }
方案2:多内置SMT组合(不推荐)
如果暂时不想写代码,也可以通过ExtractField+HoistField+ReplaceField+SetSchemaMetadata多个内置SMT组合实现,但需要为每个业务字段单独配置转换规则,后续新增字段就要修改配置,扩展性极差,仅适合字段固定且数量极少的场景。
内容的提问来源于stack exchange,提问作者Chris
相关产品推荐
相关产品推荐

