修改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

