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

如何通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 04:54:05