Kafka consumer connection status


What's inside this article ⌄
  • Kafka consumer connection state
  • Check Kafka consumer status
  • Kafka consumer connection check
  • Monitor Kafka consumer state

If you want to gather information about the interaction of Kafka Consumer with Kafka Broker, you may want to use this method:

 1public String readConsumerInfo(KafkaConsumer<String, String> consumer) {
 2    try {
 3        ConsumerNetworkClient consumerNetworkClient = getKafkaConsumerNetworkClient(consumer);
 4        KafkaClient kafkaClient = getKafkaClient(consumerNetworkClient);
 5        Object kafkaClientConnectionStates = getKafkaClientConnectionStates(kafkaClient);
 6        String result = kafkaConnectionStatesToString(kafkaClientConnectionStates);
 7        return String.format("----- %s description: %s", readConsumerClientId(consumer), result);
 8    } catch (Exception ex) {
 9        ex.printStackTrace();
10        return null;
11    }
12}

Helper methods:

 1private ConsumerNetworkClient getKafkaConsumerNetworkClient(KafkaConsumer<String, String> consumer) throws Exception {
 2    Field field = consumer.getClass().getDeclaredField("client");
 3    field.setAccessible(true);
 4    return (ConsumerNetworkClient) field.get(consumer);
 5}
 6
 7private KafkaClient getKafkaClient(ConsumerNetworkClient consumerNetworkClient) throws Exception {
 8    Field field = consumerNetworkClient.getClass().getDeclaredField("client");
 9    field.setAccessible(true);
10    return (KafkaClient) field.get(consumerNetworkClient);
11}
12
13private Object getKafkaClientConnectionStates(KafkaClient kafkaClient) throws Exception {
14    Field field = kafkaClient.getClass().getDeclaredField("connectionStates");
15    field.setAccessible(true);
16    return field.get(kafkaClient);
17}
18
19private String kafkaConnectionStatesToString(Object kafkaClusterConnectionStates) throws Exception {
20    Field field = kafkaClusterConnectionStates.getClass().getDeclaredField("nodeState");
21    field.setAccessible(true);
22    Map<String, Object> connectionStateMap = (Map<String, Object>) field.get(kafkaClusterConnectionStates);
23    StringBuilder sb = new StringBuilder();
24    for (Map.Entry<String, Object> connectionState : connectionStateMap.entrySet()) {
25        sb.append("\n").append("Id: ").append(connectionState.getKey()).append(", ").append(kafkaNodeConnectionStateToString(connectionState.getValue()));
26    }
27    return sb.toString();
28}
29
30/**
31  * You may want to read more fields from KafkaNodeConnectionState class
32  */
33private String kafkaNodeConnectionStateToString(Object kafkaNodeConnectionState) throws Exception {
34    Field stateField = kafkaNodeConnectionState.getClass().getDeclaredField("state");
35    stateField.setAccessible(true);
36    ConnectionState state = (ConnectionState) stateField.get(kafkaNodeConnectionState);
37
38    Field hostField = kafkaNodeConnectionState.getClass().getDeclaredField("host");
39    hostField.setAccessible(true);
40    String host = (String) hostField.get(kafkaNodeConnectionState);
41
42    Field authenticationExceptionField = kafkaNodeConnectionState.getClass().getDeclaredField("authenticationException");
43    authenticationExceptionField.setAccessible(true);
44    AuthenticationException authenticationException = (AuthenticationException) authenticationExceptionField.get(kafkaNodeConnectionState);
45
46    Field failedAttemptsField = kafkaNodeConnectionState.getClass().getDeclaredField("failedAttempts");
47    failedAttemptsField.setAccessible(true);
48    Long failedAttempts = (Long) failedAttemptsField.get(kafkaNodeConnectionState);
49
50    Field failedConnectAttemptsField = kafkaNodeConnectionState.getClass().getDeclaredField("failedConnectAttempts");
51    failedConnectAttemptsField.setAccessible(true);
52    Long failedConnectAttempts = (Long) failedConnectAttemptsField.get(kafkaNodeConnectionState);
53
54    return String.format("state: %s, host: %s, authenticationException: %s, failedAttempts: %s, failedConnectAttempts: %s)",
55            state, host, authenticationException, failedAttempts, failedConnectAttempts);
56}
57
58private String readConsumerClientId(KafkaConsumer<String, String> consumer) {
59    try {
60        Method method = consumer.getClass().getDeclaredMethod("getClientId");
61        method.setAccessible(true);
62        return (String) method.invoke(consumer);
63    } catch (Exception ex) {
64        ex.printStackTrace();
65        return null;
66    }
67}