如何将Debezium捕获的MySQL DATETIME转为Kafka标准时间格式
解决Debezium将MySQL DATETIME转为Unix时间戳的问题
快速全局解决方案
修改Debezium配置中的time.precision.mode参数值为connect,即可让MySQL的DATETIME列和TIMESTAMP列一样,输出为ISO 8601格式的字符串,无需修改表结构或消费端逻辑。
修改后的完整Connector配置
curl -i -X PUT -H "Content-Type:application/json" \ http://localhost:8083/connectors/mysql-debezium-test/config \ -d '{ "connector.class": "io.debezium.connector.mysql.MySqlConnector", "database.hostname": "mysql", "database.port": "3306", "database.user": "debezium", "database.password": "dbz", "database.server.id": "44", "database.server.name": "asgard2", "table.whitelist": "demo.movies,demo.second_movies", "database.history.kafka.bootstrap.servers": "broker:29092", "database.history.kafka.topic": "dbhistory.demo" , "decimal.handling.mode": "double", "include.schema.changes": "false", "snapshot.mode": "schema_only", "time.precision.mode": "connect", "transforms": "unwrap,dropTopicPrefix", "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.unwrap.drop.tombstones": "true", "transforms.unwrap.delete.handling.mode":"rewrite", "transforms.dropTopicPrefix.type":"org.apache.kafka.connect.transforms.RegexRouter", "transforms.dropTopicPrefix.regex":"asgard2.demo.(.*)", "transforms.dropTopicPrefix.replacement":"mysql2.$1", "key.converter": "org.apache.kafka.connect.json.JsonConverter", "key.converter.schemas.enable": "false", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false", "log.retention.hours": "120", "poll.interval.ms": "30000" }'
配置原理
- 当前配置中
time.precision.mode=adaptive:Debezium会根据MySQL列的类型和精度自动选择存储格式,无时区的DATETIME列默认以Unix时间戳输出,而带时区的TIMESTAMP列会转为ISO 8601字符串。 - 修改为
time.precision.mode=connect:强制所有时间类型(包括DATETIME)使用Kafka Connect的标准表示方式,统一输出为"2023-06-22T19:29:26Z"格式的字符串。
可选:针对DATETIME列单独转换
如果不想全局修改所有时间类型的处理逻辑,可以添加自定义转换器,仅对DATETIME列进行格式转换:
添加的transform配置片段
"transforms": "unwrap,dropTopicPrefix,convertDatetime", "transforms.convertDatetime.type": "io.debezium.transforms.DateTimeConverter", "transforms.convertDatetime.format": "yyyy-MM-dd'T'HH:mm:ss'Z'", "transforms.convertDatetime.field": ".*_datetime"
field参数使用正则表达式匹配目标列,示例中匹配所有以_datetime结尾的列,可根据实际列名调整规则。- 该方案仅转换指定的DATETIME列,其他时间类型保持原有处理逻辑。
内容的提问来源于stack exchange,提问作者jyablonski
相关产品推荐
相关产品推荐

