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

