Kafka Connector Timestamp转换器与ExtractNewDocumentState SMT配合失效求助
问题排查与解决方案
核心问题分析
- 转换器冲突:配置中同时启用了全局第三方时间戳转换器
timestampConverter(oryanmoshe实现)和Kafka Connect自带的TimestampConverterSMT,两者功能重叠,可能干扰彼此的执行逻辑。 - SMT参数缺失:自带的
TimestampConverterSMT默认仅处理Timestamp类型字段,但unwrap添加的ts_ms是Long类型的毫秒级时间戳,未指定source.type参数时无法识别该字段类型,导致转换失败。
修正后的配置调整
步骤1:移除冲突的全局转换器
删除全局converters相关配置项,避免与SMT冲突:
// 移除以下配置 "converters": "timestampConverter", "timestampConverter.debug": "false", "timestampConverter.format.date": "YYYY-MM-dd", "timestampConverter.format.datetime": "YYYY-MM-dd'T'HH:mm:ss'Z'", "timestampConverter.format.time": "HH:mm:ss", "timestampConverter.type": "oryanmoshe.kafka.connect.util.TimestampConverter",
步骤2:完善TimestampConverter SMT配置
为transform-name添加source.type参数,指定处理Long类型的毫秒时间戳:
"transforms.transform-name.source.type": "long"
步骤3:确认SMT执行顺序
确保transform-name(TimestampConverter)在unwrap之后执行(当前配置顺序已满足:unwrap → ReplaceField → RenameField → transform-name),保证ts_ms字段被添加后再进行转换。
最终完整配置示例
{ "connector.class": "io.debezium.connector.mongodb.MongoDbConnector", "auto.create.topics.enable": "false", "topic.creation.enable": "true", "topic.creation.default.partitions": "3", "topic.prefix": "prefix-1", "topic.creation.default.replication.factor": "3", "topic.creation.default.compression.type": "gzip", "topic.creation.default.file.delete.delay.ms": "432000000", "topic.creation.default.cleanup.policy": "delete", "topic.creation.default.retention.ms": "432000000", "tombstones.on.delete": "false", "mongodb.connection.string": "", "mongodb.name": "", "mongodb.user": "", "mongodb.password": "", "mongodb.authSource": "", "mongodb.connection.mode": "", "database.include.list": "", "name": "test-connector-01", "collection.include.list": "", "schema.history.internal.kafka.bootstrap.servers": "kafka-0.kafka-headless.kafka-connector:9092", "schema.history.internal.kafka.topic": "compare-payment-installment-mongo-schema", "transforms": "unwrap, ReplaceField, RenameField, transform-name", "transforms.ReplaceField.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", "transforms.ReplaceField.exclude": "source", "transforms.unwrap.type": "io.debezium.connector.mongodb.transforms.ExtractNewDocumentState", "transforms.unwrap.collection.expand.json.payload": "true", "transforms.unwrap.add.fields": "op,ts_ms", "transforms.unwrap.add.fields.prefix": "", "transforms.RenameField.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", "transforms.RenameField.renames": "_id:id", "transforms.transform-name.type": "org.apache.kafka.connect.transforms.TimestampConverter$Value", "transforms.transform-name.field": "ts_ms", "transforms.transform-name.format": "yyyy-MM-dd'T'HH:mm:ss.SSS'Z'", "transforms.transform-name.target.type": "string", "transforms.transform-name.source.type": "long", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false", "capture.mode": "change_streams_update_full_with_pre_image", "snapshot.mode": "never", "capture.scope": "database", "tasks.max": "1" }
验证要点
- 重启连接器后,检查Kafka主题中的消息,确认
ts_ms字段已转换为指定的yyyy-MM-dd'T'HH:mm:ss.SSS'Z'格式字符串。 - 若仍有问题,可临时开启SMT调试日志(需调整连接器日志级别),查看
TimestampConverter的执行细节,确认是否正确识别到ts_ms字段。
内容的提问来源于stack exchange,提问作者Narut Promsuparoj
相关产品推荐
相关产品推荐

