Spring Cloud Stream Kafka消费者轮询超时恢复与健康检查咨询
我们有一个基于Spring Cloud Stream的微服务,负责从Kafka Topic读取消息并写入MQTT,服务初始运行正常,但一段时间后出现以下异常,且后续无法向MQTT发布消息:
"2022-10-18 16:22:29.861 WARN 1 --- [d | tellus-mqtt] o.a.k.c.c.internals.ConsumerCoordinator : [Consumer clientId=consumer-mqtt-2, groupId=mqtt] consumer poll timeout has expired. This means the time between subsequent calls to poll() was longer than the configured max.poll.interval.ms, which typically implies that the poll loop is spending too much time processing messages. You can address this either by increasing max.poll.interval.ms or by reducing the maximum size of batches r"
我们有两个疑问:
- 是否可以通过编程方式重新订阅或从该超时中恢复?
- 能否为Actuator实现自定义健康检查以纳入消费者状态,从而让K8s自动重启异常Pod?例如配置:
management: endpoint: health: group: liveness: include: livenessstate,binders
其中binders为Kafka组件。
OutputConfig 类
@Configuration @Log4j2 @Profile("output") public class OutputConfig { private final Mqtt3ReactorClient outboundMqttClient; private final Mqtt3ReactorClient outboundRootMqttClient; private final MeterUtils meterUtils; @Autowired public OutputConfig(@Qualifier("outboundMqttClient") Mqtt3ReactorClient outboundMqttClient, @Qualifier("outboundRootMqttClient") Mqtt3ReactorClient outboundRootMqttClient, MeterUtils meterUtils) { this.outboundMqttClient = outboundMqttClient; this.outboundRootMqttClient = outboundRootMqttClient; this.meterUtils = meterUtils; log.info("Starting Output Config!"); } @Bean public Consumer<Flux<Output.GatewayNotification>> kafka() { return new Output(outboundMqttClient, meterUtils); } @Bean public Consumer<Flux<Output.GatewayNotification>> kafkaRoot() { return new Output(outboundRootMqttClient, meterUtils); } }
Output 类
@Log4j2 public class Output implements Consumer<Flux<Output.GatewayNotification>> { public static final HexFormat FORMAT = HexFormat.of().withDelimiter(" ").withUpperCase(); private final Mqtt3ReactorClient outboundMqttClient; private final MeterUtils meterUtils; public Output(Mqtt3ReactorClient outboundMqttClient, MeterUtils meterUtils) { this.outboundMqttClient = outboundMqttClient; this.meterUtils = meterUtils; } @Override public void accept(Flux<Output.GatewayNotification> gatewayNotifications) { Flux<Mqtt3Publish> messagesToPublish = gatewayNotifications .map(gatewayNotification -> Mqtt3Publish.builder() .topic(gatewayNotification.getAddress()) .qos(MqttQos.AT_LEAST_ONCE) .payload(Base64.getDecoder().decode(gatewayNotification.getPayload())) .build()); outboundMqttClient.publish(messagesToPublish) .doOnNext(publishResult -> { log.debug( "Publish acknowledged: " + FORMAT.formatHex(publishResult.getPublish().getPayloadAsBytes())); meterUtils.incrementCounter("output"); }) .doOnError(error -> log.error(error.getMessage())) .subscribe(); } @Data public static class GatewayNotification { private String address; private String payload; private Long buildingId; } }
HiveMqMqttConfig 类
@Configuration @Log4j2 public class HiveMqMqttConfig { @Value("${mqtt.endpointUrl}") private String endpointUrl; @Value("${mqtt.rootEndpointUrl}") private String rootEndpointUrl; @Value("${mqtt.inboundClientId}") private String inboundClientId; @Value("${mqtt.outboundClientId}") private String outboundClientId; @Value("${mqtt.caFilename:#{null}}") private String caFilename; @Value("${mqtt.inboundPrivateKeyFilename:#{null}}") private String inboundPrivateKeyFilename; @Value("${mqtt.inboundRootPrivateKeyFilename:#{null}}") private String inboundRootPrivateKeyFilename; @Value("${mqtt.inboundClientCertFilename:#{null}}") private String inboundClientCertFilename; @Value("${mqtt.inboundRootClientCertFilename:#{null}}") private String inboundRootClientCertFilename; @Value("${mqtt.outboundPrivateKeyFilename:#{null}}") private String outboundPrivateKeyFilename; @Value("${mqtt.outboundRootPrivateKeyFilename:#{null}}") private String outboundRootPrivateKeyFilename; @Value("${mqtt.outboundClientCertFilename:#{null}}") private String outboundClientCertFilename; @Value("${mqtt.outboundRootClientCertFilename:#{null}}") private String outboundRootClientCertFilename; @Bean(name = "inboundMqttClient") public Mqtt3ReactorClient inboundMqttClient() { var client = Mqtt3ReactorClient.from(buildMqtt3Client(endpointUrl, UUID.randomUUID().toString(), caFilename, inboundPrivateKeyFilename, inboundClientCertFilename)); connectClient(client); return client; } @Bean(name = "inboundRootMqttClient") public Mqtt3ReactorClient inboundRootMqttClient() { var client = Mqtt3ReactorClient.from(buildMqtt3Client(rootEndpointUrl, UUID.randomUUID().toString(), caFilename, inboundRootPrivateKeyFilename, inboundRootClientCertFilename)); connectClient(client); return client; } @Bean(name = "outboundMqttClient") public Mqtt3ReactorClient outboundMqttClient() { var client = Mqtt3ReactorClient.from(buildMqtt3Client(endpointUrl, UUID.randomUUID().toString(), caFilename, outboundPrivateKeyFilename, outboundClientCertFilename)); connectClient(client); return client; } @Bean(name = "outboundRootMqttClient") public Mqtt3ReactorClient outboundRootMqttClient() { var client = Mqtt3ReactorClient.from(buildMqtt3Client(rootEndpointUrl, UUID.randomUUID().toString(), caFilename, outboundRootPrivateKeyFilename, outboundRootClientCertFilename)); connectClient(client); return client; } private Mqtt3Client buildMqtt3Client(String endpointUrl, String clientId, String caFilename, String privateKeyFilename, String clientCertFilename) { log.info("Creating mqtt3 client with client id: {}", clientId); // endpoint is in the form 'protocol://host:port' String[] endpointUrlComponents = endpointUrl.split(":"); String host = endpointUrlComponents[1].substring(2); int port = Integer.parseInt(endpointUrlComponents[2]); Mqtt3ClientBuilder mqtt3ClientBuilder = Mqtt3Client.builder() .identifier(clientId) .serverHost(host) .serverPort(port) .automaticReconnectWithDefaultConfig(); try { if (caFilename != null && !caFilename.isEmpty()) { boolean isUsingKeyBasedAuthentication = privateKeyFilename != null && !privateKeyFilename.isEmpty() && clientCertFilename != null && !clientCertFilename.isEmpty(); PemFileSslContext context = isUsingKeyBasedAuthentication ? new PemFileSslContext(getStreamFromClassPathOrLocal(caFilename), getStreamFromClassPathOrLocal(privateKeyFilename), getStreamFromClassPathOrLocal(clientCertFilename)) : new PemFileSslContext(new ClassPathResource(caFilename).getInputStream()); context.getSocketFactory(); mqtt3ClientBuilder .sslConfig() .keyManagerFactory(context.getKeyManagerFactory()) .trustManagerFactory(context.getTrustManagerFactory()) .applySslConfig(); } } catch (IOException | NoSuchAlgorithmException | KeyStoreException | CertificateException | InvalidKeySpecException | UnrecoverableKeyException | PemFileSslContext.SocketFactoryCreationFailedException e) { throw new RuntimeException(e); } return mqtt3ClientBuilder.build(); } private InputStream getStreamFromClassPathOrLocal(String uri) throws IOException { return new ClassPathResource(uri).getInputStream(); } private void connectClient(Mqtt3ReactorClient mqtt3ReactorClient) { Mono<Mqtt3ConnAck> connAckSingle = mqtt3ReactorClient.connect(); connAckSingle .doOnSuccess(connAck -> log.info("Connected, " + connAck.getReturnCode())) .doOnError(throwable -> log.info("Connection failed, " + throwable.getMessage())) .subscribe(); } }
应用配置文件
management: endpoint: health: group: liveness: include: livenessstate,kafkaConsumers spring: cloud: stream: kafka: bindings: kafka-in-0: consumer: configuration: max.poll.records: 10 kafkaRoot-in-0: consumer: configuration: max.poll.records: 10 function: definition: kafka;kafkaRoot bindings: kafka-in-0: destination: output group: mqtt consumer: concurrency: 1 kafkaRoot-in-0: destination: output group: mqtt-root consumer: concurrency: 1 ... (certs/endpoints omitted)
1. 编程方式恢复Kafka消费者超时
根本问题定位
当前代码中Output类的accept方法手动调用subscribe(),导致反应式流脱离Spring Cloud Stream的上下文管理。Kafka的poll线程会被阻塞,无法在max.poll.interval.ms时限内完成下一次poll,最终触发超时并被踢出消费组。
修复与恢复方式
优先解决阻塞问题
移除手动subscribe(),让Spring Cloud Stream接管反应式流的订阅和背压处理,确保poll线程不会被阻塞:
@Override public void accept(Flux<Output.GatewayNotification> gatewayNotifications) { Flux<Mqtt3Publish> messagesToPublish = gatewayNotifications .map(gatewayNotification -> Mqtt3Publish.builder() .topic(gatewayNotification.getAddress()) .qos(MqttQos.AT_LEAST_ONCE) .payload(Base64.getDecoder().decode(gatewayNotification.getPayload())) .build()); // 移除手动subscribe(),让框架管理流的生命周期 outboundMqttClient.publish(messagesToPublish) .doOnNext(publishResult -> { log.debug( "Publish acknowledged: " + FORMAT.formatHex(publishResult.getPublish().getPayloadAsBytes())); meterUtils.incrementCounter("output"); }) .doOnError(error -> log.error("MQTT发布失败: {}", error.getMessage())) .subscribe(); // 若必须手动订阅,需确保使用非阻塞线程,优先推荐让Spring管理 }
兜底编程式重启绑定
如果需要在异常发生后强制重启Kafka消费者,可以通过BindingService实现:
@Autowired private BindingService bindingService; public void restartKafkaConsumers() { // 停止现有消费者绑定 bindingService.unbindConsumers("kafka-in-0"); bindingService.unbindConsumers("kafkaRoot-in-0"); // 重新绑定消费者 bindingService.bindConsumers(); }
此方法仅作为兜底,优先解决处理逻辑阻塞的核心问题。
2. 自定义Actuator健康检查让K8s重启Pod
方式一:利用现有组件
Spring Cloud Stream自带binders健康指示器,开启后可直接纳入liveness检查:
management: endpoint: health: group: liveness: include: livenessstate,binders health: binders: enabled: true
当Kafka消费者超时或断开连接时,binders状态会变为DOWN,K8s的liveness探针检测到后会自动重启Pod。
方式二:自定义健康检查
若需要更细粒度的状态检查(如MQTT连接、Kafka消费lag),可实现自定义HealthIndicator:
@Component public class StreamHealthIndicator implements HealthIndicator { private final Mqtt3ReactorClient outboundMqttClient; @Autowired public StreamHealthIndicator(@Qualifier("outboundMqttClient") Mqtt3ReactorClient outboundMqttClient) { this.outboundMqttClient = outboundMqttClient; } @Override public Health health() { Health.Builder builder = Health.up(); // 检查MQTT连接状态 try { outboundMqttClient.ping().block(Duration.ofSeconds(5)); builder.withDetail("mqtt-outbound", "connected"); } catch (Exception e) { builder.down().withDetail("mqtt-outbound-error", e.getMessage()); } // 可扩展添加Kafka消费者状态检查(如消费lag、连接状态) return builder.build(); } }
然后将自定义指示器加入liveness组:
management: endpoint: health: group: liveness: include: livenessstate,binders,streamHealthIndicator
内容的提问来源于stack exchange,提问作者Rod McCutcheon

