【发布时间】:2022-11-10 16:53:46
【问题描述】:
我正在尝试从 kafka 主题中获取数据,然后将其插入到 mysql 数据库表中。主题 (smartdevdbserver1.signup_db.users) 派生自另一个名为 users 的 mysql 数据库表列,并使用 debezium CDC mysql 源连接器填充。 如果有人能帮助我找出接收器连接器抛出以下错误的原因,我将不胜感激:
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}}
以下是 kafka 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": "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"
}
}
要插入数据的表模式:
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 以便接收器连接器可以创建自己的表(只是这样我看看错误是否会消失)但它创建的表只有一个字段(这是主键字段:id),并且当然,没有错误。所以我猜测接收器连接器将所有其他字段(可能)视为空值。
【问题讨论】:
标签: apache-kafka apache-kafka-connect debezium ksqldb