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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 17:00:06