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

Confluent Kafka Sink Connector插入MySQL报错:email字段无默认值

问题描述

尝试将Debezium CDC MySQL Source Connector同步的Kafka Topic(smartdevdbserver1.signup_db.users)数据插入到目标MySQL表时,Sink Connector抛出字段无默认值的错误,且自动建表仅生成主键id字段。

错误日志

connect           | java.sql.SQLException: Field 'email' doesn't have a default value
connect           | 
connect           |     at io.confluent.connect.jdbc.sink.JdbcSinkTask.put(JdbcSinkTask.java:93)
connect           |     at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:581)
connect           |     at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:333)
connect           |     at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:234)
connect           |     at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:203)
connect           |     at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:188)
connect           |     at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:243)
connect           |     at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
connect           |     at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
connect           |     at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
connect           |     at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
connect           |     at java.base/java.lang.Thread.run(Thread.java:829)
connect           | Caused by: java.sql.SQLException: java.sql.BatchUpdateException: Field 'email' doesn't have a default value       
connect           | java.sql.SQLException: Field 'email' doesn't have a default value

Kafka Topic Schema与Payload

{"schema":{"type":"struct","fields":[{"type":"int32","optional":false,"field":"id"},{"type":"string","optional":false,"field":"email"},{"type":"string","optional":false,"field":"password"},{"type":"string","optional":false,"name":"io.debezium.data.Enum","version":1,"parameters":{"allowed":"ACTIVE,INACTIVE"},"default":"INACTIVE","field":"User_status"},{"type":"string","optional":true,"field":"auth_token"}],"optional":false,"name":"smartdevdbserver1.signup_db.users.Value"},"payload":{"id":6,"email":"testing6@firstclicklimited.com","password":"$2a$10$PRGfCpjCCKqSKSf89m5M6uSRWzjlZTG7RuuJgR5MrVY.nh0BKA7Nq","User_status":"INACTIVE","auth_token":null}}

Sink Connector配置

{
    "name": "resetpassword-sink-connector",
    "config": {
        "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
        "tasks.max": "1",
        "key.converter": "org.apache.kafka.connect.json.JsonConverter",
        "value.converter": "org.apache.kafka.connect.json.JsonConverter",
        "key.converter.schemas.enable": "true",
        "value.converter.schemas.enable": "true",
        "topics": "smartdevdbserver1.signup_db.users",
        "connection.url": "jdbc:mysql://RPWD_mysql:3306/rpwd_db?user=rpwd_user&password=*xxxxxxxx*",
        "fields.whitelist": "rpwd_db.users.email,rpwd_db.users.password,rpwd_db.users.User_status,rpwd_db.users.auth_token",
        "transforms.unwrap.drop.tombstones": "false",
        "insert.mode": "upsert",
        "delete.enabled": "true",
        "table.name.format": "rpwd_db.users",
        "pk.fields": "id",
        "pk.mode": "record_key"
    }
}

目标MySQL表结构

DROP TABLE IF EXISTS `users`;
CREATE TABLE IF NOT EXISTS `users` (
  `id` int NOT NULL AUTO_INCREMENT,
  `email` varchar(255) NOT NULL,
  `password` varchar(255) NOT NULL,
  `User_status` enum('ACTIVE','INACTIVE') NOT NULL DEFAULT 'INACTIVE',
  `auth_token` varchar(255) DEFAULT NULL,
  PRIMARY KEY (`id`),
  UNIQUE KEY (`email`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

已尝试操作

开启auto.create让Connector自动建表,结果仅生成主键id字段且无报错,推测Connector未识别到其他字段。


问题分析与解决方案

核心问题

  1. fields.whitelist配置错误:该参数需指定Kafka消息(value部分)中的字段名,而非库.表.字段格式。用户当前写法导致Connector无法匹配到任何字段,插入时缺失email等必填字段。
  2. 缺少Debezium消息unwrap转换:Debezium同步的消息是包含before/after/source的Envelope格式,未配置转换的话,Sink无法提取到实际的用户数据(payload中的after部分),这也是自动建表仅生成id的原因(id来自record_key)。

修复步骤

1. 修正fields.whitelist配置

改为直接指定Kafka消息中的字段名:

"fields.whitelist": "email,password,User_status,auth_token"

2. 添加完整的unwrap转换配置

配置ExtractNewRecordState转换提取Debezium消息的实际数据:

"transforms": "unwrap",
"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
"transforms.unwrap.drop.tombstones": "false",
"transforms.unwrap.delete.handling.mode": "rewrite"

(delete.handling.mode=rewrite配合delete.enabled=true,将删除事件转为UPSERT操作处理)

3. 验证主键配置

当前pk.mode=record_key需确保Kafka消息的key包含id字段(Debezium默认将源表主键作为record_key),若key结构不符,可改为pk.mode=record_value并保留pk.fields=id,让Connector从value中提取主键。

完整修复后的Sink配置

{
    "name": "resetpassword-sink-connector",
    "config": {
        "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
        "tasks.max": "1",
        "key.converter": "org.apache.kafka.connect.json.JsonConverter",
        "value.converter": "org.apache.kafka.connect.json.JsonConverter",
        "key.converter.schemas.enable": "true",
        "value.converter.schemas.enable": "true",
        "topics": "smartdevdbserver1.signup_db.users",
        "connection.url": "jdbc:mysql://RPWD_mysql:3306/rpwd_db?user=rpwd_user&password=*xxxxxxxx*",
        "fields.whitelist": "email,password,User_status,auth_token",
        "transforms": "unwrap",
        "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
        "transforms.unwrap.drop.tombstones": "false",
        "transforms.unwrap.delete.handling.mode": "rewrite",
        "insert.mode": "upsert",
        "delete.enabled": "true",
        "table.name.format": "rpwd_db.users",
        "pk.fields": "id",
        "pk.mode": "record_key"
    }
}

额外验证

  • 查看Kafka消息key结构是否包含id:
kafka-console-consumer.sh --bootstrap-server <kafka-broker>:9092 --topic smartdevdbserver1.signup_db.users --from-beginning --property print.key=true --property key.deserializer=org.apache.kafka.connect.json.JsonConverter --property key.deserializer.schemas.enable=true
  • 若自动建表仍有问题,可手动删除目标表,开启auto.create=true,让Connector基于正确的消息结构重建表,验证字段完整性。

内容的提问来源于stack exchange,提问作者eedideyahoocom

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 16:25:36