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

同一MariaDB实例部署多Debezium连接器遇连接重置故障

问题:多Debezium连接器对应单个MariaDB Schema时周期性Connection reset报错

我们在多个MariaDB Schema上各部署了2个Debezium连接器,初期运行正常,但每隔1-2周会随机出现连接器报错,提示binlog处理时发生Connection reset,导致连接器停止,必须手动重启。仅当每个Schema只部署1个连接器时无此问题,且MariaDB服务器未宕机(同服务器上的其他连接器不受影响)。

报错日志

2022-10-31 06:18:55,106 ERROR  MySQL|scheme_1|binlog  Error during binlog processing. Last offset stored = {transaction_id=null, ts_sec=1667155787, file=mysql-bin.075628, pos=104509320, server_id=1, event=32}, binlog reader near position = mysql-bin.075628/300573885   [io.debezium.connector.mysql.MySqlStreamingChangeEventSource]
2022-10-31 06:18:55,107 ERROR  MySQL|scheme_1|binlog  Producer failure   [io.debezium.pipeline.ErrorHandler]
io.debezium.DebeziumException: Connection reset
        at io.debezium.connector.mysql.MySqlStreamingChangeEventSource.wrap(MySqlStreamingChangeEventSource.java:1189)
        at io.debezium.connector.mysql.MySqlStreamingChangeEventSource$ReaderThreadLifecycleListener.onCommunicationFailure(MySqlStreamingChangeEventSource.java:1234)
        at com.github.shyiko.mysql.binlog.BinaryLogClient.listenForEventPackets(BinaryLogClient.java:980)
        at com.github.shyiko.mysql.binlog.BinaryLogClient.connect(BinaryLogClient.java:599)
        at com.github.shyiko.mysql.binlog.BinaryLogClient$7.run(BinaryLogClient.java:857)
        at java.base/java.lang.Thread.run(Thread.java:829)
Caused by: java.net.SocketException: Connection reset
        at java.base/java.net.SocketInputStream.read(SocketInputStream.java:186)
        at java.base/java.net.SocketInputStream.read(SocketInputStream.java:140)
        at com.github.shyiko.mysql.binlog.io.BufferedSocketInputStream.read(BufferedSocketInputStream.java:59)
        at com.github.shyiko.mysql.binlog.io.ByteArrayInputStream.readWithinBlockBoundaries(ByteArrayInputStream.java:261)
        at com.github.shyiko.mysql.binlog.io.ByteArrayInputStream.read(ByteArrayInputStream.java:245)
        at com.github.shyiko.mysql.binlog.io.ByteArrayInputStream.fill(ByteArrayInputStream.java:112)
        at com.github.shyiko.mysql.binlog.io.ByteArrayInputStream.read(ByteArrayInputStream.java:105)
        at com.github.shyiko.mysql.binlog.BinaryLogClient.readPacketSplitInChunks(BinaryLogClient.java:995)
        at com.github.shyiko.mysql.binlog.BinaryLogClient.listenForEventPackets(BinaryLogClient.java:953)
        ... 3 more
2022-10-31 06:18:55,113 INFO   MySQL|scheme_1|binlog  Stopped reading binlog after 0 events, last recorded offset: {transaction_id=null, ts_sec=1667155787, file=mysql-bin.075628, pos=104509320, server_id=1, event=32}   [io.debezium.connector.mysql.MySqlStreamingChangeEventSource]
2022-10-31 06:18:55,123 ERROR  ||  WorkerSourceTask{id=scheme_1-connector-1666100046785939106-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]
org.apache.kafka.connect.errors.ConnectException: An exception occurred in the change event producer. This connector will be stopped.
        at io.debezium.pipeline.ErrorHandler.setProducerThrowable(ErrorHandler.java:50)
        at io.debezium.connector.mysql.MySqlStreamingChangeEventSource$ReaderThreadLifecycleListener.onCommunicationFailure(MySqlStreamingChangeEventSource.java:1234)
        at com.github.shyiko.mysql.binlog.BinaryLogClient.listenForEventPackets(BinaryLogClient.java:980)
        at com.github.shyiko.mysql.binlog.BinaryLogClient.connect(BinaryLogClient.java:599)
        at com.github.shyiko.mysql.binlog.BinaryLogClient$7.run(BinaryLogClient.java:857)
        at java.base/java.lang.Thread.run(Thread.java:829)
Caused by: io.debezium.DebeziumException: Connection reset
        at io.debezium.connector.mysql.MySqlStreamingChangeEventSource.wrap(MySqlStreamingChangeEventSource.java:1189)
        ... 5 more
Caused by: java.net.SocketException: Connection reset
        at java.base/java.net.SocketInputStream.read(SocketInputStream.java:186)
        at java.base/java.net.SocketInputStream.read(SocketInputStream.java:140)
        at com.github.shyiko.mysql.binlog.io.BufferedSocketInputStream.read(BufferedSocketInputStream.java:59)
        at com.github.shyiko.mysql.binlog.io.ByteArrayInputStream.readWithinBlockBoundaries(ByteArrayInputStream.java:261)
        at com.github.shyiko.mysql.binlog.io.ByteArrayInputStream.read(ByteArrayInputStream.java:245)
        at com.github.shyiko.mysql.binlog.io.ByteArrayInputStream.fill(ByteArrayInputStream.java:112)
        at com.github.shyiko.mysql.binlog.io.ByteArrayInputStream.read(ByteArrayInputStream.java:105)
        at com.github.shyiko.mysql.binlog.BinaryLogClient.readPacketSplitInChunks(BinaryLogClient.java:995)
        at com.github.shyiko.mysql.binlog.BinaryLogClient.listenForEventPackets(BinaryLogClient.java:953)
        ... 3 more
2022-10-31 06:18:55,132 INFO   ||  Stopping down connector   [io.debezium.connector.common.BaseSourceTask]

问题分析

  • MariaDB连接限制触发:同一Schema下多个连接器使用相同账号并发连接时,可能触达MariaDB的连接数上限或会话资源限制,服务器主动重置空闲或资源占用过高的连接。
  • 心跳机制缺失/不足:Debezium默认心跳间隔较长,多个并发连接处于空闲状态时,容易被MariaDB的wait_timeout或中间网络设备(如防火墙)判定为无效连接并断开。
  • Binlog会话冲突:多个连接器同时订阅同一Schema的binlog流,底层binlog客户端在请求数据、更新位点时可能产生竞争,导致服务器端连接异常。

解决方案

1. 调整MariaDB连接参数

  • 修改my.cnf配置文件,增大空闲连接超时时间:
    wait_timeout = 86400
    interactive_timeout = 86400
    
  • 检查并提升max_connections参数,确保能支撑所有连接器的并发连接需求:
    max_connections = 200
    
    重启MariaDB生效。

2. 配置Debezium连接器心跳

在连接器配置中添加以下参数,维持连接存活:

# 每30秒发送一次心跳包
heartbeat.interval.ms=30000
# 自定义心跳查询(需提前创建debezium_heartbeat表)
heartbeat.action.query="INSERT INTO debezium_heartbeat (ts) VALUES (NOW()) ON DUPLICATE KEY UPDATE ts=NOW()"

心跳表创建语句:

CREATE TABLE debezium_heartbeat (
  id INT PRIMARY KEY DEFAULT 1,
  ts TIMESTAMP NOT NULL
) ENGINE=InnoDB;

3. 优化连接器部署策略

  • 优先使用单个连接器+路由规则处理同一Schema的分流需求,避免多连接器重复监听binlog:通过transforms配置将不同表的变更路由到不同Topic,替代多连接器部署。
  • 若必须部署多个连接器,为每个连接器分配独立的MariaDB账号,避免同一账号的并发连接触发服务器限制。

4. 升级Debezium版本

升级到Debezium 2.0及以上稳定版本,新版本修复了多个binlog客户端并发连接时的会话管理问题,提升了连接稳定性。

5. 启用自动恢复机制

在Kafka Connect的连接器配置中添加错误容忍参数,让连接器在报错后自动尝试恢复:

errors.tolerance=all
errors.deadletterqueue.topic.name=dlq-debezium-scheme-1
errors.deadletterqueue.context.headers.enable=true

内容的提问来源于stack exchange,提问作者Kęstutis Keršis

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 18:15:40