如何通过SMT反序列化Avro schema并在写入ES Sink Connector前移除schema
解决方案
关于Elasticsearch sink是否需要schema
Elasticsearch sink connector原生支持无schema模式,只要你在连接器配置中设置schema.ignore=true,它就不会读取Connect记录的schema,直接将记录值作为纯JSON结构序列化写入ES,不需要你维护schema逻辑。
你之前报错的核心原因是:直接返回schema为null的记录,但记录值仍为绑定了原schema的Struct类型,Kafka Connect会校验值类型和schema的匹配性,自然会抛出不兼容错误。
自定义SMT的实现方案
你不需要实现makeUpdatedSchema逻辑,按以下步骤调整即可:
- 第一步:在SMT的处理逻辑中,先将传入的带schema的
Struct值递归转换为普通Java集合类型:结构化对象转HashMap、数组转List、基础数据类型保持原样,完全剥离和原schema的绑定。 - 第二步:转换完成后,再调用
newRecord(record, null, 转换后的纯集合对象)返回新记录,此时返回的就是合法的无schema记录。 - 第三步:调整连接器配置,补充以下参数:
# ES sink忽略schema schema.ignore=true # 若使用JsonConverter,关闭schema封装 value.converter.schemas.enable=false
调整完成后,整个流程就不需要再处理schema重建逻辑,既简化了代码,也能避免每条记录重构schema的性能开销。
内容的提问来源于stack exchange,提问作者RyanQuey
相关产品推荐
相关产品推荐

