更新Debezium连接器table.include.list后服务异常终止问题
问题:更新Debezium连接器table.include.list后出现Leader not known错误
在Kubernetes集群中部署了Debezium与Confluent Schema Registry,用于流转数据库事件,二者直接连接Confluent平台内的Kafka集群。每次更新连接器的table.include.list配置后,Debezium都会停止运行,抛出如下错误,同时Schema Registry也出现相关报错。
Debezium错误日志
org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:223) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:149) at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.convertTransformedRecord(AbstractWorkerSourceTask.java:474) at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.sendRecords(AbstractWorkerSourceTask.java:387) at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.execute(AbstractWorkerSourceTask.java:354) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:189) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:244) at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.run(AbstractWorkerSourceTask.java:72) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: org.apache.kafka.connect.errors.DataException: Failed to serialize Avro data from topic XXXXXXX : at io.confluent.connect.avro.AvroConverter.fromConnectData(AvroConverter.java:93) at org.apache.kafka.connect.storage.Converter.fromConnectData(Converter.java:64) at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.lambda$convertTransformedRecord$5(AbstractWorkerSourceTask.java:474) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:173) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:207) ... 12 more Caused by: org.apache.kafka.common.errors.SerializationException: Error registering Avro schema{"type":"record","name":"Key","namespace":"XXXXX","fields":[{XXXXX}],"connect.name":"XXXXXXX"} at io.confluent.kafka.serializers.AbstractKafkaSchemaSerDe.toKafkaException(AbstractKafkaSchemaSerDe.java:259) at io.confluent.kafka.serializers.AbstractKafkaAvroSerializer.serializeImpl(AbstractKafkaAvroSerializer.java:156) at io.confluent.connect.avro.AvroConverter$Serializer.serialize(AvroConverter.java:153) at io.confluent.connect.avro.AvroConverter.fromConnectData(AvroConverter.java:86) ... 16 more Caused by: io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException: Leader not known. io.confluent.kafka.schemaregistry.rest.exceptions.RestUnknownLeaderException: Leader not known.
Schema Registry错误日志
io.confluent.kafka.schemaregistry.rest.exceptions.RestUnknownLeaderException: Leader not known. at io.confluent.kafka.schemaregistry.rest.exceptions.Errors.unknownLeaderException(Errors.java:178) at io.confluent.kafka.schemaregistry.rest.resources.SubjectVersionsResource.register(SubjectVersionsResource.java:304) at jdk.internal.reflect.GeneratedMethodAccessor19.invoke(Unknown Source) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:566) at org.glassfish.jersey.server.model.internal.ResourceMethodInvocationHandlerFactory.lambda$static$0(ResourceMethodInvocationHandlerFactory.java:52) at org.glassfish.jersey.server.model.internal.AbstractJavaResourceMethodDispatcher$1.run(AbstractJavaResourceMethodDispatcher.java:124) at org.glassfish.jersey.server.model.internal.AbstractJavaResourceMethodDispatcher.invoke(AbstractJavaResourceMethodDispatcher.java:167) at org.glassfish.jersey.server.model.internal.JavaResourceMethodDispatcherProvider$VoidOutInvoker.doDispatch(JavaResourceMethodDispatcherProvider.java:159) at org.glassfish.jersey.server.model.internal.AbstractJavaResourceMethodDispatcher.dispatch(AbstractJavaResourceMethodDispatcher.java:79) at org.glassfish.jersey.server.model.ResourceMethodInvoker.invoke(ResourceMethodInvoker.java:475) at org.glassfish.jersey.server.model.ResourceMethodInvoker.apply(ResourceMethodInvoker.java:397) at org.glassfish.jersey.server.model.ResourceMethodInvoker.apply(ResourceMethodInvoker.java:81) at org.glassfish.jersey.server.ServerRuntime$1.run(ServerRuntime.java:255) at org.glassfish.jersey.internal.Errors$1.call(Errors.java:248) at org.glassfish.jersey.internal.Errors$1.call(Errors.java:244) at org.glassfish.jersey.internal.Errors.process(Errors.java:292) at org.glassfish.jersey.internal.Errors.process(Errors.java:274) at org.glassfish.jersey.internal.Errors.process(Errors.java:244) at org.glassfish.jersey.process.internal.RequestScope.runInScope(RequestScope.java:265) at org.glassfish.jersey.server.ServerRuntime.process(ServerRuntime.java:234) at org.glassfish.jersey.server.ApplicationHandler.handle(ApplicationHandler.java:680) at org.glassfish.jersey.servlet.WebComponent.serviceImpl(WebComponent.java:394) at org.glassfish.jersey.servlet.ServletContainer.serviceImpl(ServletContainer.java:386) at org.glassfish.jersey.servlet.ServletContainer.doFilter(ServletContainer.java:561) at org.glassfish.jersey.servlet.ServletContainer.doFilter(ServletContainer.java:502) at org.glassfish.jersey.servlet.ServletContainer.doFilter(ServletContainer.java:439) at org.eclipse.jetty.servlet.FilterHolder.doFilter(FilterHolder.java:193) at org.eclipse.jetty.servlet.ServletHandler$Chain.doFilter(ServletHandler.java:1601) at org.eclipse.jetty.servlet.ServletHandler.doHandle(ServletHandler.java:548) at org.eclipse.jetty.server.handler.ScopedHandler.nextHandle(ScopedHandler.java:233) at org.eclipse.jetty.server.session.SessionHandler.doHandle(SessionHandler.java:1624) at org.eclipse.jetty.server.handler.ScopedHandler.nextHandle(ScopedHandler.java:233) at org.eclipse.jetty.server.handler.ContextHandler.doHandle(ContextHandler.java:1434) at org.eclipse.jetty.server.handler.ScopedHandler.nextScope(ScopedHandler.java:188) at org.eclipse.jetty.servlet.ServletHandler.doScope(ServletHandler.java:501) at org.eclipse.jetty.server.session.SessionHandler.doScope(SessionHandler.java:1594) at org.eclipse.jetty.server.handler.ScopedHandler.nextScope(ScopedHandler.java:186) at org.eclipse.jetty.server.handler.ContextHandler.doScope(ContextHandler.java:1349) at org.eclipse.jetty.server.handler.ScopedHandler.handle(ScopedHandler.java:141) at org.eclipse.jetty.server.handler.HandlerCollection.handle(HandlerCollection.java:146) at org.eclipse.jetty.server.handler.HandlerCollection.handle(HandlerCollection.java:146) at org.eclipse.jetty.server.handler.StatisticsHandler.handle(StatisticsHandler.java:179) at org.eclipse.jetty.server.handler.ContextHandlerCollection.handle(ContextHandlerCollection.java:234) at org.eclipse.jetty.server.handler.gzip.GzipHandler.handle(GzipHandler.java:763) at org.eclipse.jetty.server.handler.HandlerWrapper.handle(HandlerWrapper.java:127) at org.eclipse.jetty.server.Server.handle(Server.java:516) at org.eclipse.jetty.server.HttpChannel.lambda$handle$1(HttpChannel.java:388) at org.eclipse.jetty.server.HttpChannel.dispatch(HttpChannel.java:633) at org.eclipse.jetty.server.HttpChannel.handle(HttpChannel.java:380) at org.eclipse.jetty.server.HttpConnection.onFillable(HttpConnection.java:277) at org.eclipse.jetty.io.AbstractConnection$ReadCallback.succeeded(AbstractConnection.java:311) at org.eclipse.jetty.io.FillInterest.fillable(FillInterest.java:105) at org.eclipse.jetty.io.ChannelEndPoint$1.run(ChannelEndPoint.java:104) at org.eclipse.jetty.util.thread.strategy.EatWhatYouKill.runTask(EatWhatYouKill.java:338) at org.eclipse.jetty.util.thread.strategy.EatWhatYouKill.doProduce(EatWhatYouKill.java:315) at org.eclipse.jetty.util.thread.strategy.EatWhatYouKill.tryProduce(EatWhatYouKill.java:173) at org.eclipse.jetty.util.thread.strategy.EatWhatYouKill.run(EatWhatYouKill.java:131) at org.eclipse.jetty.util.thread.ReservedThreadExecutor$ReservedThread.run(ReservedThreadExecutor.java:386) at org.eclipse.jetty.util.thread.QueuedThreadPool.runJob(QueuedThreadPool.java:883) at org.eclipse.jetty.util.thread.QueuedThreadPool$Runner.run(QueuedThreadPool.java:1034) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: io.confluent.kafka.schemaregistry.exceptions.UnknownLeaderException: Register schema request failed since leader is unknown at io.confluent.kafka.schemaregistry.storage.KafkaSchemaRegistry.registerOrForward(KafkaSchemaRegistry.java:610) at io.confluent.kafka.schemaregistry.rest.resources.SubjectVersionsResource.register(SubjectVersionsResource.java:284) ... 60 more
应用版本信息
- Debezium版本:2.1.2.Final
- Confluent版本:7.0.4
- Kafka版本:2.8
Schema Registry配置
SCHEMA_REGISTRY_LISTENERS: http://0.0.0.0:8081 SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: xxxxx SCHEMA_REGISTRY_KAFKASTORE_TOPIC: _schemas_xxxxx
排查与解决方案
排查方向
- Schema Registry存储主题异常
核心原因是Schema Registry无法找到其存储主题(_schemas_xxxxx)的Leader节点。更新table.include.list后,Debezium为新增表生成新Avro Schema并尝试注册,此时Schema Registry需写入自身存储主题,但Leader缺失导致失败。 - Kafka集群连通性
- 验证
SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS配置的Kafka地址是否正确,确保Schema Registry能访问Broker节点。 - 在Kubernetes的Schema Registry Pod内执行
nc -zv <kafka-broker> <port>测试网络连通性。
- 验证
- 存储主题状态检查
登录Kafka集群执行以下命令查看主题状态:kafka-topics.sh --describe --topic _schemas_xxxxx --bootstrap-server <kafka-bootstrap>- 检查
Leader列是否为有效值,-1表示Leader缺失。 - 确认副本数和ISR(同步副本)数量,确保ISR中有可用节点。
- 检查
- Schema注册并发压力
一次性新增大量表会触发大量Schema注册请求,可能导致Schema Registry或Kafka负载过高,引发Leader选举异常。
解决方案
- 修复存储主题Leader问题
- 手动触发Leader选举:
kafka-preferred-replica-election.sh --bootstrap-server <kafka-bootstrap> --topic _schemas_xxxxx - 若副本数不足(如为1),调整副本数提升可用性:
kafka-topics.sh --alter --topic _schemas_xxxxx --replication-factor 3 --bootstrap-server <kafka-bootstrap>
- 手动触发Leader选举:
- 修正Kafka连接配置
- 确保
SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS配置的是Kafka集群节点列表或内部服务名,而非单个节点,避免单点故障。 - 若Kafka开启安全认证(如SASL/SSL),需添加对应认证参数到Schema Registry配置中。
- 确保
- 优化Debezium更新方式
- 分批次更新
table.include.list,减少并发Schema注册请求。 - 临时调整Debezium错误容忍配置(生产环境谨慎使用):设置
errors.tolerance=all,避免单个Schema注册失败导致连接器停止。
- 分批次更新
- 提升Schema Registry可用性
- 部署多实例Schema Registry,通过负载均衡器提供服务,避免单点故障。
- 确保所有Schema Registry实例配置一致,尤其是存储主题和Kafka连接参数。
内容的提问来源于stack exchange,提问作者Jonathan Chevalier
相关产品推荐
相关产品推荐

