Kafka Connect Standalone启动挂起问题排查求助
Kafka Connect JDBC连接器启动后挂起排查方案
问题描述
使用Kafka Connect搭配confluentinc/kafka-connect-jdbc连接器及MySQL驱动捕获MySQL数据变更,执行connect-standalone.sh启动后进程挂起,日志输出到org.apache.kafka.connect.runtime.ConnectorConfig$EnrichedConnectorConfig:376后无后续连接器工作日志,仅间隔一段时间出现Jetty会话清理的DEBUG日志。
启动日志
[2023-03-14 14:43:49,701] INFO [jdbc_source_mysql_03|worker] Attempting to open connection #1 to MySql (io.confluent.connect.jdbc.util.CachedConnectionProvider:80) [2023-03-14 14:43:49,702] INFO [jdbc_source_mysql_03|worker] --------------- START GETTING CONNECTION (io.confluent.connect.jdbc.dialect.GenericDatabaseDialect:250) [2023-03-14 14:43:50,118] INFO [jdbc_source_mysql_03|worker] --------------- STOP GETTING CONNECTION (io.confluent.connect.jdbc.dialect.GenericDatabaseDialect:252) [2023-03-14 14:43:50,121] INFO [jdbc_source_mysql_03|worker] Starting thread to monitor tables. (io.confluent.connect.jdbc.source.TableMonitorThread:82) [2023-03-14 14:43:50,123] DEBUG [jdbc_source_mysql_03|worker] Using MySql dialect to check support for [TABLE] (io.confluent.connect.jdbc.dialect.GenericDatabaseDialect:481) [2023-03-14 14:43:50,124] INFO SourceConnectorConfig values: config.action.reload = restart connector.class = io.confluent.connect.jdbc.JdbcSourceConnector errors.log.enable = true errors.log.include.messages = false errors.retry.delay.max.ms = 60000 errors.retry.timeout = 0 errors.tolerance = none header.converter = null key.converter = null name = jdbc_source_mysql_03 predicates = [] tasks.max = 1 topic.creation.groups = [] transforms = [] value.converter = null (org.apache.kafka.connect.runtime.SourceConnectorConfig:376) [2023-03-14 14:43:50,125] INFO EnrichedConnectorConfig values: config.action.reload = restart connector.class = io.confluent.connect.jdbc.JdbcSourceConnector errors.log.enable = true errors.log.include.messages = false errors.retry.delay.max.ms = 60000 errors.retry.timeout = 0 errors.tolerance = none header.converter = null key.converter = null name = jdbc_source_mysql_03 predicates = [] tasks.max = 1 topic.creation.groups = [] transforms = [] value.converter = null (org.apache.kafka.connect.runtime.ConnectorConfig$EnrichedConnectorConfig:376) [2023-03-14 14:43:50,132] DEBUG [jdbc_source_mysql_03|worker] Used MySql dialect to find table types: [TABLE] (io.confluent.connect.jdbc.dialect.GenericDatabaseDialect:500) [2023-03-14 14:43:50,132] DEBUG [jdbc_source_mysql_03|worker] Using MySql dialect to get TABLE (io.confluent.connect.jdbc.dialect.GenericDatabaseDialect:429) [2023-03-14 14:43:50,156] DEBUG [jdbc_source_mysql_03|worker] Used MySql dialect to find 1 TABLE (io.confluent.connect.jdbc.dialect.GenericDatabaseDialect:442) [2023-03-14 14:43:50,157] DEBUG [jdbc_source_mysql_03|worker] Got the following tables: ["mytestdb"."test11"] (io.confluent.connect.jdbc.source.TableMonitorThread:176) [2023-03-14 14:43:50,157] INFO [jdbc_source_mysql_03|worker] --------------- START GETTING CONNECTION (io.confluent.connect.jdbc.dialect.GenericDatabaseDialect:250) [2023-03-14 14:43:50,183] INFO [jdbc_source_mysql_03|worker] --------------- STOP GETTING CONNECTION (io.confluent.connect.jdbc.dialect.GenericDatabaseDialect:252) [2023-03-14 14:43:50,185] INFO [jdbc_source_mysql_03|worker] SourceConnectorConfig values: config.action.reload = restart connector.class = io.confluent.connect.jdbc.JdbcSourceConnector errors.log.enable = true errors.log.include.messages = false errors.retry.delay.max.ms = 60000 errors.retry.timeout = 0 errors.tolerance = none header.converter = null key.converter = null name = jdbc_source_mysql_03 predicates = [] tasks.max = 1 topic.creation.groups = [] transforms = [] value.converter = null (org.apache.kafka.connect.runtime.SourceConnectorConfig:376) [2023-03-14 14:43:50,189] INFO [jdbc_source_mysql_03|worker] EnrichedConnectorConfig values: config.action.reload = restart connector.class = io.confluent.connect.jdbc.JdbcSourceConnector errors.log.enable = true errors.log.include.messages = false errors.retry.delay.max.ms = 60000 errors.retry.timeout = 0 errors.tolerance = none header.converter = null key.converter = null name = jdbc_source_mysql_03 predicates = [] tasks.max = 1 topic.creation.groups = [] transforms = [] value.converter = null (org.apache.kafka.connect.runtime.ConnectorConfig$EnrichedConnectorConfig:376) [2023-03-14 14:43:50,187] DEBUG [jdbc_source_mysql_03|worker] Based on the supplied filtering rules, the tables available to read from include: `mytestdb`.`test11` (io.confluent.connect.jdbc.source.TableMonitorThread:127) [2023-03-14 14:43:50,190] DEBUG [jdbc_source_mysql_03|worker] Based on the supplied filtering rules, the tables available to read from include: `mytestdb`.`test11` (io.confluent.connect.jdbc.source.TableMonitorThread:127) [2023-03-14 14:53:48,918] DEBUG node0 scavenging sessions (org.eclipse.jetty.server.session:241) [2023-03-14 14:53:48,918] DEBUG org.eclipse.jetty.server.session.SessionHandler2025928493==dftMaxIdleSec=-1 scavenging sessions (org.eclipse.jetty.server.session:1251) [2023-03-14 14:53:48,919] DEBUG org.eclipse.jetty.server.session.SessionHandler2025928493==dftMaxIdleSec=-1 scavenging session ids [] (org.eclipse.jetty.server.session:1259) [2023-03-14 14:53:48,919] DEBUG org.eclipse.jetty.server.session.DefaultSessionCache@44bd4b0a[evict=-1,removeUnloadable=false,saveOnCreate=false,saveOnInactiveEvict=false] checking expiration on [] (org.eclipse.jetty.server.session:688)
配置文件
connect-standalone.properties
bootstrap.servers=172.22.10.21:9092,172.22.10.22:9092,172.22.10.23:9092 key.converter=org.apache.kafka.connect.json.JsonConverter value.converter=org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable=true value.converter.schemas.enable=true offset.storage.file.filename=/tmp/connect.offsets offset.flush.interval.ms=10000 plugin.path=/app/plugin
source-connector.properties
name=jdbc_source_mysql_03 connector.class=io.confluent.connect.jdbc.JdbcSourceConnector connection.url=jdbc:mysql://172.22.10.24:3306/mytestdb connection.user=mytest connection.password=password numeric.mapping=best_fit tasks.max=1 mode=incrementing #timestamp.column.name=modifiedAt incrementing.column.name=id table.whitelist=mytestdb.test11 table.types=TABLE timestamp.initial=-1 topic.prefix=mysql1.mytestdb. table.poll.interval.ms=5000 errors.log.enable=true
排查原因及解决方案
从日志看,连接器已成功连接MySQL并识别到目标表mytestdb.test11,但后续任务初始化或数据读取环节出现阻塞,以下是可能的原因及解决方法:
1. 日志级别过低隐藏关键信息
当前日志仅输出INFO和部分DEBUG,任务执行时的异常或阻塞细节未被记录。
- 解决方法:修改Connect日志配置,将
io.confluent.connect.jdbc和org.apache.kafka.connect.runtime的日志级别调整为DEBUG,重启后查看详细日志定位阻塞点。
2. Kafka集群连接异常
虽配置了bootstrap.servers,但Connect可能无法正常与Kafka集群建立连接,导致任务因无法发送数据阻塞。
- 解决方法:
- 在Connect所在机器执行
telnet 172.22.10.21 9092测试Kafka端口连通性; - 检查Kafka集群节点状态,确认服务正常运行;
- 开启DEBUG日志,查看是否有Kafka连接超时或认证失败信息。
- 在Connect所在机器执行
3. 偏移量文件权限问题
偏移量文件/tmp/connect.offsets权限异常,导致Connect无法读写偏移量,阻塞任务初始化。
- 解决方法:
- 执行
ls -l /tmp/connect.offsets检查文件权限,确保Connect进程拥有读写权限; - 若文件不存在,确认
/tmp目录权限正常,Connect可创建文件; - 权限异常时执行
chown <connect-user>:<connect-group> /tmp/connect.offsets修改权限。
- 执行
4. JDBC连接器与MySQL驱动版本不兼容
kafka-connect-jdbc版本与MySQL驱动版本不匹配,导致连接后静默阻塞。
- 解决方法:
- 确认版本兼容性(如Connector 10.x搭配MySQL 8.0驱动,Connector 5.x搭配MySQL 5.x驱动);
- 替换为兼容的驱动版本后重启Connect。
5. 目标表数据量过大导致首次加载阻塞
test11表数据量过大,首次全量读取占用大量资源,导致进程看似挂起。
- 解决方法:
- 查看系统CPU、内存、磁盘IO使用情况,确认是否因资源耗尽阻塞;
- 先测试小表验证连接器功能,再针对大表调整
batch.max.rows参数控制每次读取行数,避免一次性加载过多数据。
6. MySQL会话阻塞
MySQL侧会话被锁等待等情况阻塞,导致Connect查询无法返回结果。
- 解决方法:
- 在MySQL中执行
SHOW PROCESSLIST;查看是否有长时间运行的查询或锁等待; - 清理阻塞会话或等待锁释放,观察Connect是否恢复。
- 在MySQL中执行
内容的提问来源于stack exchange,提问作者Pakorn K
相关产品推荐
相关产品推荐

