ConsumerGroupService.java 9.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219
  1. package com.provectus.kafka.ui.service;
  2. import com.provectus.kafka.ui.model.ConsumerGroupOrderingDTO;
  3. import com.provectus.kafka.ui.model.InternalConsumerGroup;
  4. import com.provectus.kafka.ui.model.KafkaCluster;
  5. import com.provectus.kafka.ui.model.SortOrderDTO;
  6. import java.util.ArrayList;
  7. import java.util.Comparator;
  8. import java.util.HashMap;
  9. import java.util.List;
  10. import java.util.Map;
  11. import java.util.Properties;
  12. import java.util.UUID;
  13. import java.util.function.ToIntFunction;
  14. import java.util.stream.Collectors;
  15. import javax.annotation.Nullable;
  16. import lombok.RequiredArgsConstructor;
  17. import lombok.Value;
  18. import org.apache.commons.lang3.StringUtils;
  19. import org.apache.kafka.clients.admin.ConsumerGroupDescription;
  20. import org.apache.kafka.clients.admin.OffsetSpec;
  21. import org.apache.kafka.clients.consumer.ConsumerConfig;
  22. import org.apache.kafka.clients.consumer.KafkaConsumer;
  23. import org.apache.kafka.common.TopicPartition;
  24. import org.apache.kafka.common.serialization.BytesDeserializer;
  25. import org.apache.kafka.common.utils.Bytes;
  26. import org.springframework.stereotype.Service;
  27. import reactor.core.publisher.Flux;
  28. import reactor.core.publisher.Mono;
  29. import reactor.util.function.Tuple2;
  30. import reactor.util.function.Tuples;
  31. @Service
  32. @RequiredArgsConstructor
  33. public class ConsumerGroupService {
  34. private final AdminClientService adminClientService;
  35. private Mono<List<InternalConsumerGroup>> getConsumerGroups(
  36. ReactiveAdminClient ac,
  37. List<ConsumerGroupDescription> descriptions) {
  38. return Flux.fromIterable(descriptions)
  39. // 1. getting committed offsets for all groups
  40. .flatMap(desc -> ac.listConsumerGroupOffsets(desc.groupId())
  41. .map(offsets -> Tuples.of(desc, offsets)))
  42. .collectMap(Tuple2::getT1, Tuple2::getT2)
  43. .flatMap((Map<ConsumerGroupDescription, Map<TopicPartition, Long>> groupOffsetsMap) -> {
  44. var tpsFromGroupOffsets = groupOffsetsMap.values().stream()
  45. .flatMap(v -> v.keySet().stream())
  46. .collect(Collectors.toSet());
  47. // 2. getting end offsets for partitions with in committed offsets
  48. return ac.listOffsets(tpsFromGroupOffsets, OffsetSpec.latest())
  49. .map(endOffsets ->
  50. descriptions.stream()
  51. .map(desc -> {
  52. var groupOffsets = groupOffsetsMap.get(desc);
  53. var endOffsetsForGroup = new HashMap<>(endOffsets);
  54. endOffsetsForGroup.keySet().retainAll(groupOffsets.keySet());
  55. // 3. gathering description & offsets
  56. return InternalConsumerGroup.create(desc, groupOffsets, endOffsetsForGroup);
  57. })
  58. .collect(Collectors.toList()));
  59. });
  60. }
  61. @Deprecated // need to migrate to pagination
  62. public Mono<List<InternalConsumerGroup>> getAllConsumerGroups(KafkaCluster cluster) {
  63. return adminClientService.get(cluster)
  64. .flatMap(ac -> describeConsumerGroups(ac, null)
  65. .flatMap(descriptions -> getConsumerGroups(ac, descriptions)));
  66. }
  67. public Mono<List<InternalConsumerGroup>> getConsumerGroupsForTopic(KafkaCluster cluster,
  68. String topic) {
  69. return adminClientService.get(cluster)
  70. // 1. getting topic's end offsets
  71. .flatMap(ac -> ac.listOffsets(topic, OffsetSpec.latest())
  72. .flatMap(endOffsets -> {
  73. var tps = new ArrayList<>(endOffsets.keySet());
  74. // 2. getting all consumer groups
  75. return ac.listConsumerGroups()
  76. .flatMap((List<String> groups) ->
  77. Flux.fromIterable(groups)
  78. // 3. for each group trying to find committed offsets for topic
  79. .flatMap(g ->
  80. ac.listConsumerGroupOffsets(g, tps)
  81. .map(offsets -> Tuples.of(g, offsets)))
  82. .filter(t -> !t.getT2().isEmpty())
  83. .collectMap(Tuple2::getT1, Tuple2::getT2)
  84. )
  85. .flatMap((Map<String, Map<TopicPartition, Long>> groupOffsets) ->
  86. // 4. getting description for groups with non-emtpy offsets
  87. ac.describeConsumerGroups(new ArrayList<>(groupOffsets.keySet()))
  88. .map((Map<String, ConsumerGroupDescription> descriptions) ->
  89. descriptions.values().stream().map(desc ->
  90. // 5. gathering and filter non-target-topic data
  91. InternalConsumerGroup.create(
  92. desc, groupOffsets.get(desc.groupId()), endOffsets)
  93. .retainDataForPartitions(p -> p.topic().equals(topic))
  94. )
  95. .collect(Collectors.toList())));
  96. }));
  97. }
  98. @Value
  99. public static class ConsumerGroupsPage {
  100. List<InternalConsumerGroup> consumerGroups;
  101. int totalPages;
  102. }
  103. public Mono<ConsumerGroupsPage> getConsumerGroupsPage(
  104. KafkaCluster cluster,
  105. int page,
  106. int perPage,
  107. @Nullable String search,
  108. ConsumerGroupOrderingDTO orderBy,
  109. SortOrderDTO sortOrderDto
  110. ) {
  111. var comparator = sortOrderDto.equals(SortOrderDTO.ASC)
  112. ? getPaginationComparator(orderBy)
  113. : getPaginationComparator(orderBy).reversed();
  114. return adminClientService.get(cluster).flatMap(ac ->
  115. describeConsumerGroups(ac, search).flatMap(descriptions ->
  116. getConsumerGroups(
  117. ac,
  118. descriptions.stream()
  119. .sorted(comparator)
  120. .skip((long) (page - 1) * perPage)
  121. .limit(perPage)
  122. .collect(Collectors.toList())
  123. ).map(cgs -> new ConsumerGroupsPage(
  124. cgs,
  125. (descriptions.size() / perPage) + (descriptions.size() % perPage == 0 ? 0 : 1))))
  126. );
  127. }
  128. private Comparator<ConsumerGroupDescription> getPaginationComparator(ConsumerGroupOrderingDTO
  129. orderBy) {
  130. switch (orderBy) {
  131. case NAME:
  132. return Comparator.comparing(ConsumerGroupDescription::groupId);
  133. case STATE:
  134. ToIntFunction<ConsumerGroupDescription> statesPriorities = cg -> {
  135. switch (cg.state()) {
  136. case STABLE:
  137. return 0;
  138. case COMPLETING_REBALANCE:
  139. return 1;
  140. case PREPARING_REBALANCE:
  141. return 2;
  142. case EMPTY:
  143. return 3;
  144. case DEAD:
  145. return 4;
  146. case UNKNOWN:
  147. return 5;
  148. default:
  149. return 100;
  150. }
  151. };
  152. return Comparator.comparingInt(statesPriorities);
  153. case MEMBERS:
  154. return Comparator.comparingInt(cg -> -cg.members().size());
  155. default:
  156. throw new IllegalStateException("Unsupported order by: " + orderBy);
  157. }
  158. }
  159. private Mono<List<ConsumerGroupDescription>> describeConsumerGroups(ReactiveAdminClient ac,
  160. @Nullable String search) {
  161. return ac.listConsumerGroups()
  162. .map(groupIds -> groupIds
  163. .stream()
  164. .filter(groupId -> search == null || StringUtils.containsIgnoreCase(groupId, search))
  165. .collect(Collectors.toList()))
  166. .flatMap(ac::describeConsumerGroups)
  167. .map(cgs -> new ArrayList<>(cgs.values()));
  168. }
  169. public Mono<InternalConsumerGroup> getConsumerGroupDetail(KafkaCluster cluster,
  170. String consumerGroupId) {
  171. return adminClientService.get(cluster)
  172. .flatMap(ac -> ac.describeConsumerGroups(List.of(consumerGroupId))
  173. .filter(m -> m.containsKey(consumerGroupId))
  174. .map(r -> r.get(consumerGroupId))
  175. .flatMap(descr ->
  176. getConsumerGroups(ac, List.of(descr))
  177. .filter(groups -> !groups.isEmpty())
  178. .map(groups -> groups.get(0))));
  179. }
  180. public Mono<Void> deleteConsumerGroupById(KafkaCluster cluster,
  181. String groupId) {
  182. return adminClientService.get(cluster)
  183. .flatMap(adminClient -> adminClient.deleteConsumerGroups(List.of(groupId)));
  184. }
  185. public KafkaConsumer<Bytes, Bytes> createConsumer(KafkaCluster cluster) {
  186. return createConsumer(cluster, Map.of());
  187. }
  188. public KafkaConsumer<Bytes, Bytes> createConsumer(KafkaCluster cluster,
  189. Map<String, Object> properties) {
  190. Properties props = new Properties();
  191. props.putAll(cluster.getProperties());
  192. props.put(ConsumerConfig.CLIENT_ID_CONFIG, "kafka-ui-" + UUID.randomUUID());
  193. props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, cluster.getBootstrapServers());
  194. props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, BytesDeserializer.class);
  195. props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, BytesDeserializer.class);
  196. props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
  197. props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
  198. props.put(ConsumerConfig.ALLOW_AUTO_CREATE_TOPICS_CONFIG, "false");
  199. props.putAll(properties);
  200. return new KafkaConsumer<>(props);
  201. }
  202. }