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}