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

修改Kafka Connect配置后为何重发全部消息?

问题原因与解决办法

核心问题1:自定义Query未添加增量过滤条件

你当前配置使用了query: "SELECT * FROM outbox"而非table.name参数,而Kafka Connect的JdbcSourceConnector不会自动为自定义Query添加timestamp+incrementing模式的过滤条件。这意味着每次连接器轮询时,都会执行全表查询,自然会重放所有消息。

即使你配置了timestamp.initial=-1,这个参数也无法作用于自定义Query——它只对连接器自动生成的查询(使用table.name时)生效,用来设置初始的时间戳过滤阈值。

核心问题2:模式切换导致Offset不兼容

从incrementing模式切换到timestamp+incrementing模式时,两种模式的Offset存储结构完全不同:

  • 旧模式仅记录incrementing列的最大值(如{"incrementing": 100})
  • 新模式需要同时记录timestamp和incrementing列的阈值(如{"timestamp": 1690000000000, "incrementing": 100})

连接器启动时找不到对应新模式的Offset,会判定为首次运行,触发初始查询——而由于你的自定义Query没有过滤逻辑,直接拉取了全表数据。


解决办法

方案1:改用table.name自动生成查询

将query参数替换为table.name,让连接器自动生成包含timestamp和incrementing过滤条件的查询语句,同时timestamp.delay.interval.ms也会自动生效:

"mode": "timestamp+incrementing",
"incrementing.column.name": "id",
"timestamp.column.name": "created",
"timestamp.delay.interval.ms": 10000,
"timestamp.initial": -1,
"db.timezone": "Europe/Vienna",
"table.name": "outbox"

方案2:修改自定义Query,手动添加过滤条件

如果必须保留自定义Query,需要在语句中手动添加timestamp和incrementing的过滤逻辑,且占位符顺序要和模式匹配(先timestamp,再incrementing):

"query": "SELECT * FROM outbox WHERE created >= ? AND id >= ?",
"mode": "timestamp+incrementing",
"incrementing.column.name": "id",
"timestamp.column.name": "created",
"timestamp.delay.interval.ms": 10000,
"timestamp.initial": -1,
"db.timezone": "Europe/Vienna"

这样连接器会自动将Offset中的timestamp和id值填充到占位符中,实现增量同步,同时timestamp.delay.interval.ms会控制每次轮询只拉取created <= (当前时间 - 10000ms)的记录。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 19:32:46