Envoy Sidecar就绪前Kafka消费启动引发异常的解决方法
问题描述
Kubernetes生产环境部署应用时出现大量错误,Kafka消费者抛出序列化异常,根源是无法连接Kafka Schema Registry(连接被拒绝)。排查发现,Kafka消费者在Envoy Sidecar就绪前就开始消费消息,导致无法正常访问依赖的Schema Registry。需要延迟Kafka消息消费直到Envoy Sidecar就绪,且无法使用Kubernetes就绪探针,寻求有效解决经验。
错误栈信息
15:12:03.495 [xx-0-C-1] ERROR o.s.k.l.KafkaMessageListenerContainer - Consumer exception java.lang.IllegalStateException: This error handler cannot process 'SerializationException's directly; please consider configuring an 'ErrorHandlingDeserializer' in the value and/or key deserializer at org.springframework.kafka.listener.DefaultErrorHandler.handleOtherException(DefaultErrorHandler.java:192) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.handleConsumerException(KafkaMessageListenerContainer.java:1925) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1348) at java.base/java.util.concurrent.CompletableFuture$AsyncRun.run(CompletableFuture.java:1804) at java.base/java.lang.Thread.run(Thread.java:840) Caused by: org.apache.kafka.common.errors.RecordDeserializationException: Error deserializing key/value for partition xxxxxxxx-15 at offset xxxxx. If needed, please seek past the record to continue consumption. at org.apache.kafka.clients.consumer.internals.CompletedFetch.parseRecord(CompletedFetch.java:331) at org.apache.kafka.clients.consumer.internals.CompletedFetch.fetchRecords(CompletedFetch.java:283) at org.apache.kafka.clients.consumer.internals.FetchCollector.fetchRecords(FetchCollector.java:168) at org.apache.kafka.clients.consumer.internals.FetchCollector.collectFetch(FetchCollector.java:134) at org.apache.kafka.clients.consumer.internals.Fetcher.collectFetch(Fetcher.java:145) at org.apache.kafka.clients.consumer.internals.LegacyKafkaConsumer.pollForFetches(LegacyKafkaConsumer.java:666) at org.apache.kafka.clients.consumer.internals.LegacyKafkaConsumer.poll(LegacyKafkaConsumer.java:617) at org.apache.kafka.clients.consumer.internals.LegacyKafkaConsumer.poll(LegacyKafkaConsumer.java:590) at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:874) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollConsumer(KafkaMessageListenerContainer.java:1625) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doPoll(KafkaMessageListenerContainer.java:1600) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1405) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1296) ... 2 common frames omitted Caused by: org.apache.kafka.common.errors.SerializationException: Error retrieving Avro value schema for id 19222 at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer$DeserializationContext.schemaFromRegistry(AbstractKafkaAvroDeserializer.java:345) at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer.deserialize(AbstractKafkaAvroDeserializer.java:152) at io.confluent.kafka.serializers.KafkaAvroDeserializer.deserialize(KafkaAvroDeserializer.java:53) at org.apache.kafka.common.serialization.Deserializer.deserialize(Deserializer.java:62) at org.apache.kafka.common.serialization.Deserializer.deserialize(Deserializer.java:73) at org.apache.kafka.clients.consumer.internals.CompletedFetch.parseRecord(CompletedFetch.java:321) ... 14 common frames omitted Caused by: java.net.ConnectException: Connection refused at java.base/sun.nio.ch.Net.pollConnect(Native Method) at java.base/sun.nio.ch.Net.pollConnectNow(Net.java:672) at java.base/sun.nio.ch.NioSocketImpl.timedFinishConnect(NioSocketImpl.java:547) at java.base/sun.nio.ch.NioSocketImpl.connect(NioSocketImpl.java:602) at java.base/java.net.Socket.connect(Socket.java:639) at java.base/sun.net.NetworkClient.doConnect(NetworkClient.java:178) at java.base/sun.net.www.http.HttpClient.openServer(HttpClient.java:533) at java.base/sun.net.www.http.HttpClient.openServer(HttpClient.java:638) at java.base/sun.net.www.http.HttpClient.<init>(HttpClient.java:281) at java.base/sun.net.www.http.HttpClient.New(HttpClient.java:386) at java.base/sun.net.www.http.HttpClient.New(HttpClient.java:408) at java.base/sun.net.www.protocol.http.HttpURLConnection.getNewHttpClient(HttpURLConnection.java:1312) at java.base/sun.net.www.protocol.http.HttpURLConnection.plainConnect0(HttpURLConnection.java:1245) at java.base/sun.net.www.protocol.http.HttpURLConnection.plainConnect(HttpURLConnection.java:1131) at java.base/sun.net.www.protocol.http.HttpURLConnection.connect(HttpURLConnection.java:1060) at java.base/sun.net.www.protocol.http.HttpURLConnection.getInputStream0(HttpURLConnection.java:1690) at java.base/sun.net.www.protocol.http.HttpURLConnection.getInputStream(HttpURLConnection.java:1614) at java.base/java.net.HttpURLConnection.getResponseCode(HttpURLConnection.java:529) at io.confluent.kafka.schemaregistry.client.rest.RestService.sendHttpRequest(RestService.java:281) at io.confluent.kafka.schemaregistry.client.rest.RestService.httpRequest(RestService.java:371) at io.confluent.kafka.schemaregistry.client.rest.RestService.getId(RestService.java:840) at io.confluent.kafka.schemaregistry.client.rest.RestService.getId(RestService.java:813) at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.getSchemaByIdFromRegistry(CachedSchemaRegistryClient.java:294) at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.getSchemaBySubjectAndId(CachedSchemaRegistryClient.java:417) at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer$DeserializationContext.schemaFromRegistry(AbstractKafkaAvroDeserializer.java:342) ... 19 common frames omitted
可行解决方案
1. 应用启动阶段添加Envoy就绪检测
在应用启动流程中加入循环检测逻辑,直到能通过Envoy正常访问Schema Registry再继续初始化:
@Component public class EnvoyReadyChecker implements ApplicationRunner { @Value("${spring.kafka.properties.schema.registry.url}") private String schemaRegistryUrl; private static final Logger log = LoggerFactory.getLogger(EnvoyReadyChecker.class); @Override public void run(ApplicationArguments args) throws Exception { int maxAttempts = 30; int waitSeconds = 2; int attempts = 0; boolean isReady = false; while (attempts < maxAttempts && !isReady) { try { URL url = new URL(schemaRegistryUrl + "/subjects"); HttpURLConnection connection = (HttpURLConnection) url.openConnection(); connection.setRequestMethod("GET"); int responseCode = connection.getResponseCode(); if (responseCode >= 200 && responseCode < 300) { isReady = true; log.info("Envoy就绪,Schema Registry可访问"); } else { log.warn("Schema Registry返回非2xx状态码:{}", responseCode); } } catch (IOException e) { log.warn("连接Schema Registry失败(Envoy未就绪):{}", e.getMessage()); } if (!isReady) { Thread.sleep(waitSeconds * 1000); attempts++; } } if (!isReady) { throw new IllegalStateException("Envoy超时未就绪"); } } }
2. 延迟Kafka消费者容器启动
针对Spring Kafka,自定义容器工厂,在启动消费者容器前先执行就绪检测:
@Bean public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory( ConsumerFactory<String, Object> consumerFactory) { ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); factory.setContainerCustomizer(container -> { try { checkEnvoyReady(); container.start(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("Envoy就绪检测中断", e); } }); return factory; } private void checkEnvoyReady() throws InterruptedException { // 复用上述就绪检测逻辑 }
3. 启动脚本中检测Envoy状态
如果使用Istio,直接在应用启动脚本中检测Envoy admin端口(默认15000)的就绪状态:
# 等待Envoy就绪 until curl -s http://localhost:15000/server_info | grep -q "READY"; do echo "等待Envoy Sidecar就绪..." sleep 2 done echo "Envoy就绪,启动应用..." # 启动应用 java -jar app.jar
4. 配置Kafka消费者错误重试
即使消费者提前启动,也可以通过错误重试机制应对初期的连接异常:
- 添加
ErrorHandlingDeserializer配置:
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer spring.kafka.consumer.properties.spring.deserializer.value.delegate.class=io.confluent.kafka.serializers.KafkaAvroDeserializer
- 配置带重试的错误处理器:
@Bean public DefaultErrorHandler errorHandler(KafkaTemplate<String, Object> kafkaTemplate) { return new DefaultErrorHandler( new DeadLetterPublishingRecoverer(kafkaTemplate), new FixedBackOff(2000L, 10)); // 重试10次,每次间隔2秒 }
内容的提问来源于stack exchange,提问作者Ajay Kumar
相关产品推荐
相关产品推荐

