Kafka MySQL JDBC连接器使用WHERE子句查询失败问题
Kafka MySQL JDBC源连接器WHERE子句查询报错问题
问题场景
使用Kafka的MySQL JDBC源连接器时,自定义带WHERE子句的查询无法正常运行,触发SQL语法错误。
连接器配置文件
{ "name": "mysql-jdbc", "config": { "connector.class" : "io.confluent.connect.jdbc.JdbcSourceConnector", "connection.url" : "jdbc:mysql://mysqldb:3306/silicon", "connection.user" : "root", "connection.password" : "root", "mode" : "incrementing", "incrementing.column.name": "id", "query": "select * from (select * from credit_lines where id=2) test;", "topic.prefix" : "JDBC.test_db_test", "validate.non.null" : "false", "poll.interval.ms" : "1000" } }
日志中的查询配置信息
connect | (org.apache.kafka.connect.runtime.tracing.TracerConfig) connect | [2023-03-01 19:42:30,132] INFO Initializing: org.apache.kafka.connect.runtime.TransformationChain{} (org.apache.kafka.connect.runtime.Worker) connect | [2023-03-01 19:42:30,133] INFO [Worker clientId=connect-1, groupId=jdbc_source_connector] Finished starting connectors and tasks (org.apache.kafka.connect.runtime.distributed.DistributedHerder) connect | [2023-03-01 19:42:30,133] INFO Starting JDBC source task (io.confluent.connect.jdbc.source.JdbcSourceTask) connect | [2023-03-01 19:42:30,133] INFO JdbcSourceTaskConfig values: connect | batch.max.rows = 100 connect | catalog.pattern = null connect | connection.attempts = 3 connect | connection.backoff.ms = 10000 connect | connection.password = [hidden] connect | connection.url = jdbc:mysql://mysqldb:3306/silicon connect | connection.user = root connect | db.timezone = UTC connect | dialect.name = connect | incrementing.column.name = id connect | mode = incrementing connect | numeric.mapping = null connect | numeric.precision.mapping = false connect | poll.interval.ms = 1000 connect | query = select * from (select * from credit_lines where id=2) test; connect | query.retry.attempts = -1 connect | query.suffix = connect | quote.sql.identifiers = ALWAYS connect | schema.pattern = null connect | table.blacklist = [] connect | table.monitoring.startup.polling.limit.ms = 10000 connect | table.poll.interval.ms = 60000 connect | table.types = [TABLE] connect | table.whitelist = [] connect | tables = [] connect | timestamp.column.name = [] connect | timestamp.delay.interval.ms = 0 connect | timestamp.granularity = connect_logical connect | timestamp.initial = null connect | topic.prefix = JDBC.test_db_test connect | transaction.isolation.mode = DEFAULT connect | validate.non.null = false connect | (io.confluent.connect.jdbc.source.JdbcSourceTaskConfig)
报错信息
ERROR WorkerSourceTask{id=mysql-jdbc-0} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted (org.apache.kafka.connect.runtime.WorkerTask) connect | org.apache.kafka.connect.errors.ConnectException: java.sql.SQLSyntaxErrorException: You have an error in your SQL syntax; check the manual that corresponds to your MySQL server version for the right syntax to use near 'WHERE `id` > -1 ORDER BY `id` ASC' at line 1 connect | at io.confluent.connect.jdbc.source.JdbcSourceTask.poll(JdbcSourceTask.java:452) connect | at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.poll(AbstractWorkerSourceTask.java:470) connect | at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.execute(AbstractWorkerSourceTask.java:349) connect | at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:201) connect | at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:256) connect | at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.run(AbstractWorkerSourceTask.java:75) 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.SQLSyntaxErrorException: You have an error in your SQL syntax; check the manual that corresponds to your MySQL server version for the right syntax to use near 'WHERE `id` > -1 ORDER BY `id` ASC' at line 1
问题原因
当配置了mode: incrementing和incrementing.column.name: id时,JDBC源连接器会自动在自定义查询语句末尾追加增量条件WHERE id > [last_value] ORDER BY id ASC。但你的自定义查询末尾加了分号,导致拼接后的SQL变成:
select * from (select * from credit_lines where id=2) test; WHERE `id` > -1 ORDER BY `id` ASC
分号会终止当前SQL语句,后续的WHERE子句就成了无效语法,触发报错。
解决方案
方案1:移除查询末尾的分号
修改query配置,去掉最后的分号,让连接器能正确拼接增量条件:
"query": "select * from (select * from credit_lines where id=2) test"
方案2:切换为全量模式(不需要增量同步)
如果不需要增量同步,将mode改为bulk,连接器不会自动追加WHERE子句,每次执行完整的自定义查询:
"mode": "bulk"
方案3:调整查询逻辑适配增量模式
如果需要保留增量同步,确保查询语句返回的结果包含incrementing.column.name指定的字段,并且语句末尾无分号,让连接器可以正确拼接条件。
内容的提问来源于stack exchange,提问作者Maninder
相关产品推荐
相关产品推荐

