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

使用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表本身:

  1. UserName字段对应的Schema未正确携带泛型元数据,导致Kryo无法解析类型信息
  2. 该表的Schema发生过演化(如字段类型变更),但Debezium未正确同步最新Schema
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 07:15:34