使用Debezium同步SQLServer至Kafka时加密字段触发NPE问题
问题
使用Debezium实现SQLServer到Kafka的CDC(变更数据捕获)同步,需对指定列进行加密。环境为K8s上运行2个Kafka Connect实例,总计部署约50个SQLServer同步连接器。
某连接器配置片段如下:
{"name": "live.sql.users", ... "transforms.unwrap.delete.handling.mode": "drop", "transforms": "unwrap,cipher", "predicates.isTombstone.type": "org.apache.kafka.connect.transforms.predicates.RecordIsTombstone", "transforms.unwrap.drop.tombstones": "false", "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.cipher.predicate": "isTombstone", "transforms.cipher.negate": "true", "transforms.cipher.cipher_data_keys": "[ { \"identifier\": \"my-key\", \"material\": { \"primaryKeyId\": 1000000001, \"key\": [ { \"keyData\": { \"typeUrl\": \"type.googleapis.com/google.crypto.tink.AesGcmKey\", \"value\": \"GhDLeulEJRDC8/19NMUXqw2jK\", \"keyMaterialType\": \"SYMMETRIC\" }, \"status\": \"ENABLED\", \"keyId\": 2000000002, \"outputPrefixType\": \"TINK\" } ] } } ]", "transforms.cipher.type": "com.github.hpgrahsl.kafka.connect.transforms.kryptonite.CipherField$Value", "transforms.cipher.cipher_mode": "ENCRYPT", "predicates": "isTombstone", "transforms.cipher.field_config": "[{\"name\":\"Password\"},{\"name\":\"MobNumber\"}, {\"name\":\"UserName\"}]", "transforms.cipher.cipher_data_key_identifier": "my-key" ... }
应用配置后,调用/connectors/<connector_name>/status接口很快出现错误:
org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:206) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:132) at org.apache.kafka.connect.runtime.TransformationChain.apply(TransformationChain.java:50) at org.apache.kafka.connect.runtime.WorkerSourceTask.sendRecords(WorkerSourceTask.java:346) at org.apache.kafka.connect.runtime.WorkerSourceTask.execute(WorkerSourceTask.java:261) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:191) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:240) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: org.apache.kafka.connect.errors.DataException: error: ENCRYPT of field path 'UserName' having data 'deleted605' failed unexpectedly at com.github.hpgrahsl.kafka.connect.transforms.kryptonite.RecordHandler.processField(RecordHandler.java:90) at com.github.hpgrahsl.kafka.connect.transforms.kryptonite.SchemaawareRecordHandler.lambda$matchFields$0(SchemaawareRecordHandler.java:73) at java.base/java.util.ArrayList.forEach(ArrayList.java:1541) at java.base/java.util.Collections$UnmodifiableCollection.forEach(Collections.java:1085) at com.github.hpgrahsl.kafka.connect.transforms.kryptonite.SchemaawareRecordHandler.matchFields(SchemaawareRecordHandler.java:50) at com.github.hpgrahsl.kafka.connect.transforms.kryptonite.CipherField.processWithSchema(CipherField.java:163) at com.github.hpgrahsl.kafka.connect.transforms.kryptonite.CipherField.apply(CipherField.java:140) at org.apache.kafka.connect.runtime.PredicatedTransformation.apply(PredicatedTransformation.java:56) at org.apache.kafka.connect.runtime.TransformationChain.lambda$apply$0(TransformationChain.java:50) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:156) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:190) ... 11 more Caused by: java.lang.NullPointerException at com.esotericsoftware.kryo.util.DefaultGenerics.nextGenericTypes(DefaultGenerics.java:77) at com.esotericsoftware.kryo.serializers.FieldSerializer.pushTypeVariables(FieldSerializer.java:144) at com.esotericsoftware.kryo.serializers.FieldSerializer.write(FieldSerializer.java:102) at com.esotericsoftware.kryo.Kryo.writeObject(Kryo.java:627) at com.github.hpgrahsl.kafka.connect.transforms.kryptonite.RecordHandler.processField(RecordHandler.java:75) ... 21 more
注:相同配置在其他连接器上可正常运行,需排查该NPE问题的原因及解决方案。
原因分析
从栈轨迹看,NPE发生在Kryo序列化库的DefaultGenerics.nextGenericTypes方法中,说明在对字段值加密前的序列化环节,遇到了泛型类型信息缺失的情况。
结合场景,虽然配置与其他连接器一致,但问题出在当前同步的users表本身:
UserName字段对应的Schema未正确携带泛型元数据,导致Kryo无法解析类型信息- 该表的Schema发生过演化(如字段类型变更),但Debezium未正确同步最新Schema
- Kryptonite加密插件与当前Kryo版本存在兼容性问题,处理特定类型数据时触发NPE
解决方案
针对上述原因,可按以下步骤排查解决:
1. 验证目标字段的Schema定义
- 查看Debezium为
users表生成的Schema,确认UserName字段类型是否为STRING,是否存在nullable等特殊属性 - 对比其他正常连接器的表Schema,找出差异点(如该字段是否为新增或最近修改过类型)
2. 调整加密配置,隔离异常字段
- 临时从
field_config中移除UserName字段,观察连接器是否恢复正常,确认是否为该字段导致问题 - 若插件支持,添加
null_value_strategy配置,处理空值或特殊值场景
3. 更新加密插件或依赖版本
- 检查当前Kryptonite插件版本,升级到官方修复过类似NPE问题的最新稳定版
- 确认Kafka Connect使用的Kryo版本与插件兼容,避免版本冲突
4. 重置连接器偏移量,重新同步
- 重置该连接器的偏移量,触发全量数据同步,避免历史Schema演化遗留的异常
- 调整Debezium的Schema生成配置(如
include.schema.changes),确保Schema携带完整类型元数据
5. 切换序列化方式(临时方案)
- 若Kryptonite插件支持,将序列化方式从Kryo改为JSON,验证是否解决NPE问题
内容的提问来源于stack exchange,提问作者Karim Tawfik
相关产品推荐
相关产品推荐

