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

Kafka Connect JdbcSinkConnector插入报错:无效字节序列问题求助

解决方案:Kafka Connect同步SQL Server到PostgreSQL时的0x00字符报错问题

优先方案:原样写入包含不可打印字符的数据

PostgreSQL的TEXT类型不允许存储UTF-8无效的0x00空字符,但可以通过以下方式实现原样存储:

1. 将目标字段类型改为BYTEA

  • 停止出错的Sink连接器,删除自动创建的目标表
  • 手动重建表,将原TEXT类型的字段替换为BYTEA类型(BYTEA支持存储任意字节序列,包括0x00)
  • 修改Sink配置:
    • 关闭自动建表:"auto.create": "false"
    • 添加Cast转换,将源NVARCHAR字段转为字节类型:
      "transforms": "castToBytes",
      "transforms.castToBytes.type": "org.apache.kafka.connect.transforms.Cast$Value",
      "transforms.castToBytes.spec": "your_target_field:bytes"
      
  • 重启Sink连接器,即可将包含0x00的原数据以字节形式存入BYTEA字段。读取时需自行处理字节到字符串的转换。

备选方案:移除不可打印字符

如果不需要保留不可打印字符,可在源端或Sink端过滤:

1. 源端Debezium添加自定义SMT过滤

编写简单的Single Message Transform(SMT)移除字符串中的0x00或其他不可打印字符:

  • 实现org.apache.kafka.connect.transforms.Transformation接口,核心逻辑示例:
    public class RemoveNonPrintableChars<R extends ConnectRecord<R>> implements Transformation<R> {
        private Set<String> fieldsToProcess;
    
        @Override
        public R apply(R record) {
            Struct value = (Struct) record.value();
            for (String field : fieldsToProcess) {
                if (value.get(field) instanceof String) {
                    String original = (String) value.get(field);
                    // 移除0x00或扩展正则移除所有不可打印字符
                    String cleaned = original.replaceAll("\\x00", "");
                    value.put(field, cleaned);
                }
            }
            return record.newRecord(record.topic(), record.kafkaPartition(), record.keySchema(), record.key(), value.schema(), value, record.timestamp());
        }
    
        // 实现configure、close等方法,读取配置的需要处理的字段列表
    }
    
  • 将打包后的jar放入Kafka Connect插件目录,修改Debezium源连接器配置:
    "transforms": "cleanFields",
    "transforms.cleanFields.type": "com.yourorg.transforms.RemoveNonPrintableChars$Value",
    "transforms.cleanFields.fields": "field1,field2"
    

2. Sink端使用Script转换过滤

利用Kafka Connect内置的Script转换,通过Groovy脚本处理字符串:

  • 修改Sink连接器配置:
    "transforms": "cleanStrings",
    "transforms.cleanStrings.type": "org.apache.kafka.connect.transforms.Script$Value",
    "transforms.cleanStrings.script": "clean_non_printable.groovy"
    
  • 创建clean_non_printable.groovy脚本,内容如下:
    def value = value
    // 遍历所有字符串字段,移除0x00
    value.each { key, val ->
        if (val instanceof String) {
            value[key] = val.replaceAll("\\x00", "")
            // 如需移除所有不可打印字符,替换为:val.replaceAll("[^\\x20-\\x7E\\x0A\\x0D]", "")
        }
    }
    return value
    
  • 将脚本放在Kafka Connect可访问的路径,重启Sink连接器即可生效。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 14:01:16