Spring Cloud Stream - Polled Consumer halts with frequent rebalance message

Viewed 154

I have an application that is using multiple consumers and polled consumers to listen in on messages. The consumers seems to be working fine but the polled consumers seems to be having issues. From what I can see, every call to PollableMessageSource.poll() leads to frequent printing of Attempt to heartbeat failed since group is rebalancing until finally the entire application stalls.

Here's one of the polled consumers I've created:

@Autowired
@Qualifier("consumeEGSRequest-in-0")
private PollableMessageSource eGSRequestSource;

@Scheduled(fixedDelay = 5_000)
public void pollEGSRequest() {
    if (someCondition) {
    }
    boolean polled = eGSRequestSource.poll(m -> {
        RBWrapper rbw = (RBWrapper) m.getPayload();
        synchronized (sConfig.getSE()) {
            sConfig.getSE().add(rbw.getRequestId());
        }
        gSBProxy.executeGS(rbw.getParam(), rbw.getRequestId());  // asynchronously process message
    }, new ParameterizedTypeReference<RBWrapper>() {});
    log.debug("Poll Status :: " + polled);
}

Here's the application yml:

spring:
    cloud:
        stream:
            kafka:
                binder:
                    consumer-properties:
                        max.poll.interval.ms: 3600000
                        max.poll.records: 1
            binders:
                kafka1:
                    environment:
                        spring:
                            cloud:
                                stream:
                                    kafka:
                                        binder:
                                            brokers: localhost:9092
                                            zkNodes: localhost:2181
                                            auto-create-topics: true
                                            #min-partition-count: 2
                    type: kafka
            pollable-source: consumeEGSRequest;consumeESRRequest;consumeSPRequest;consumeSSRPRequest;consumeGCGF4CPRequest;consumeGD4CPRequest;consumeGCGF4LPRequest;consumeGD4LPRequest
            function:
                definition: consumeIDRequest;consumeDSRequest;consumeSOAPRequest;consumeLSSDRequest;consumeSOAP4ARequest;consumeDOResponse4S;consumeGMI4SResponse;consumeOS4SResponse
            bindings:
                consumeSPRequest-in-0:
                    binder: kafka1
                    destination: SP_REQ
                    group: consumer_cloud_stream1
                cr-2-sg-sM-sP-resp:
                    binder: kafka1
                    destination: SP_RESP
                    
                consumeIDRequest-in-0:
                    binder: kafka1
                    destination: ID_REQ
                    group: consumer_cloud_stream1
                cr-2-sg-sM-iD-resp:
                    binder: kafka1
                    destination: ID_RESP
                    
                consumeSSRPRequest-in-0:
                    binder: kafka1
                    destination: SSRP_REQ
                    group: consumer_cloud_stream1
                cr-2-sg-sM-sSRP-resp:
                    binder: kafka1
                    destination: SSRP_RESP
                    
                consumeEGSRequest-in-0:
                    binder: kafka1
                    destination: EGS_REQ
                    group: consumer_cloud_stream1
                cr-2-sg-gSB-eGS-resp:
                    binder: kafka1
                    destination: EGS_RESP
                    
                consumeDSRequest-in-0:
                    binder: kafka1
                    destination: DS_REQ
                    group: consumer_cloud_stream1
                cr-2-sg-nO-dS-resp:
                    binder: kafka1
                    destination: DS_RESP
                    
                consumeESRRequest-in-0:
                    binder: kafka1
                    destination: ESR_REQ
                    group: consumer_cloud_stream1
                cr-2-sg-nO-eSR-resp:
                    binder: kafka1
                    destination: ESR_RESP
                    
                consumeSOAPRequest-in-0:
                    binder: kafka1
                    destination: SOAP_REQ
                    group: consumer_cloud_stream1
                cr-2-sg-sU-sOAP-resp:
                    binder: kafka1
                    destination: SOAP_RESP
                    
                consumeLSSDRequest-in-0:
                    binder: kafka1
                    destination: LSSD_REQ
                    group: consumer_cloud_stream1
                cr-2-sg-sU-lSSD-resp:
                    binder: kafka1
                    destination: LSSD_RESP
                    
                consumeSOAP4ARequest-in-0:
                    binder: kafka1
                    destination: SOAP4A_REQ
                    group: consumer_cloud_stream1
                ac-2-sg-sU-sOAP-resp:
                    binder: kafka1
                    destination: SOAP4A_RESP
                    
                consumeGCGF4CPRequest-in-0:
                    binder: kafka1
                    destination: GCGF4CP_REQ
                    group: consumer_cloud_stream1
                cp-2-sg-sU-gCGF-resp:
                    binder: kafka1
                    destination: GCGF4CP_RESP
                    
                consumeGD4CPRequest-in-0:
                    binder: kafka1
                    destination: GD4CP_REQ
                    group: consumer_cloud_stream1
                cp-2-sg-sU-gD-resp:
                    binder: kafka1
                    destination: GD4CP_RESP
                    
                consumeGCGF4LPRequest-in-0:
                    binder: kafka1
                    destination: GCGF4LP_REQ
                    group: consumer_cloud_stream1
                lp-2-sg-sU-gCGF-resp:
                    binder: kafka1
                    destination: GCGF4LP_RESP
                    
                consumeGD4LPRequest-in-0:
                    binder: kafka1
                    destination: GD4LP_REQ
                    group: consumer_cloud_stream1
                lp-2-sg-sU-gD-resp:
                    binder: kafka1
                    destination: GD4LP_RESP
                    
                    
                sg-2-cr-dO-req:
                    binder: kafka1
                    destination: DO4S_REQ
                consumeDOResponse4S-in-0:
                    binder: kafka1
                    destination: DO4S_RESP
                    group: consumer_cloud_stream1
                    
                sg-2-cr-lFDD-gMI-req:
                    binder: kafka1
                    destination: GMI_REQ
                consumeGMI4SResponse-in-0:
                    binder: kafka1
                    destination: GMI_RESP
                    group: consumer_cloud_stream1
                    
                sg-2-cr-mSFU-oS-req:
                    binder: kafka1
                    destination: OS_REQ
                consumeOS4SResponse-in-0:
                    binder: kafka1
                    destination: OS_RESP
                    group: consumer_cloud_stream1
                
                

Here's the link to the full application log,

Application log

Spring boot Version: 2.6.4. Would appreciate the help. Thanks.

0 Answers
Related