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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 19:54:52