KafkaConnectService.java 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309
  1. package com.provectus.kafka.ui.service;
  2. import com.fasterxml.jackson.core.type.TypeReference;
  3. import com.fasterxml.jackson.databind.ObjectMapper;
  4. import com.provectus.kafka.ui.connect.api.KafkaConnectClientApi;
  5. import com.provectus.kafka.ui.connect.model.ConnectorStatus;
  6. import com.provectus.kafka.ui.connect.model.ConnectorStatusConnector;
  7. import com.provectus.kafka.ui.connect.model.ConnectorTopics;
  8. import com.provectus.kafka.ui.connect.model.TaskStatus;
  9. import com.provectus.kafka.ui.exception.NotFoundException;
  10. import com.provectus.kafka.ui.exception.ValidationException;
  11. import com.provectus.kafka.ui.mapper.ClusterMapper;
  12. import com.provectus.kafka.ui.mapper.KafkaConnectMapper;
  13. import com.provectus.kafka.ui.model.ConnectDTO;
  14. import com.provectus.kafka.ui.model.ConnectorActionDTO;
  15. import com.provectus.kafka.ui.model.ConnectorDTO;
  16. import com.provectus.kafka.ui.model.ConnectorPluginConfigValidationResponseDTO;
  17. import com.provectus.kafka.ui.model.ConnectorPluginDTO;
  18. import com.provectus.kafka.ui.model.ConnectorStateDTO;
  19. import com.provectus.kafka.ui.model.ConnectorTaskStatusDTO;
  20. import com.provectus.kafka.ui.model.FullConnectorInfoDTO;
  21. import com.provectus.kafka.ui.model.KafkaCluster;
  22. import com.provectus.kafka.ui.model.NewConnectorDTO;
  23. import com.provectus.kafka.ui.model.TaskDTO;
  24. import com.provectus.kafka.ui.model.connect.InternalConnectInfo;
  25. import com.provectus.kafka.ui.util.ReactiveFailover;
  26. import java.util.List;
  27. import java.util.Map;
  28. import java.util.Optional;
  29. import java.util.function.Predicate;
  30. import java.util.stream.Collectors;
  31. import java.util.stream.Stream;
  32. import javax.annotation.Nullable;
  33. import lombok.RequiredArgsConstructor;
  34. import lombok.SneakyThrows;
  35. import lombok.extern.slf4j.Slf4j;
  36. import org.apache.commons.lang3.StringUtils;
  37. import org.springframework.stereotype.Service;
  38. import org.springframework.web.reactive.function.client.WebClientResponseException;
  39. import reactor.core.publisher.Flux;
  40. import reactor.core.publisher.Mono;
  41. import reactor.util.function.Tuples;
  42. @Service
  43. @Slf4j
  44. @RequiredArgsConstructor
  45. public class KafkaConnectService {
  46. private final ClusterMapper clusterMapper;
  47. private final KafkaConnectMapper kafkaConnectMapper;
  48. private final ObjectMapper objectMapper;
  49. private final KafkaConfigSanitizer kafkaConfigSanitizer;
  50. public Flux<ConnectDTO> getConnects(KafkaCluster cluster) {
  51. return Flux.fromIterable(
  52. Optional.ofNullable(cluster.getOriginalProperties().getKafkaConnect())
  53. .map(lst -> lst.stream().map(clusterMapper::toKafkaConnect).toList())
  54. .orElse(List.of())
  55. );
  56. }
  57. public Flux<FullConnectorInfoDTO> getAllConnectors(final KafkaCluster cluster,
  58. @Nullable final String search) {
  59. return getConnects(cluster)
  60. .flatMap(connect -> getConnectorNames(cluster, connect.getName()).map(cn -> Tuples.of(connect.getName(), cn)))
  61. .flatMap(pair -> getConnector(cluster, pair.getT1(), pair.getT2()))
  62. .flatMap(connector ->
  63. getConnectorConfig(cluster, connector.getConnect(), connector.getName())
  64. .map(config -> InternalConnectInfo.builder()
  65. .connector(connector)
  66. .config(config)
  67. .build()
  68. )
  69. )
  70. .flatMap(connectInfo -> {
  71. ConnectorDTO connector = connectInfo.getConnector();
  72. return getConnectorTasks(cluster, connector.getConnect(), connector.getName())
  73. .collectList()
  74. .map(tasks -> InternalConnectInfo.builder()
  75. .connector(connector)
  76. .config(connectInfo.getConfig())
  77. .tasks(tasks)
  78. .build()
  79. );
  80. })
  81. .flatMap(connectInfo -> {
  82. ConnectorDTO connector = connectInfo.getConnector();
  83. return getConnectorTopics(cluster, connector.getConnect(), connector.getName())
  84. .map(ct -> InternalConnectInfo.builder()
  85. .connector(connector)
  86. .config(connectInfo.getConfig())
  87. .tasks(connectInfo.getTasks())
  88. .topics(ct.getTopics())
  89. .build()
  90. );
  91. })
  92. .map(kafkaConnectMapper::fullConnectorInfoFromTuple)
  93. .filter(matchesSearchTerm(search));
  94. }
  95. private Predicate<FullConnectorInfoDTO> matchesSearchTerm(@Nullable final String search) {
  96. if (search == null) {
  97. return c -> true;
  98. }
  99. return connector -> getStringsForSearch(connector)
  100. .anyMatch(string -> StringUtils.containsIgnoreCase(string, search));
  101. }
  102. private Stream<String> getStringsForSearch(FullConnectorInfoDTO fullConnectorInfo) {
  103. return Stream.of(
  104. fullConnectorInfo.getName(),
  105. fullConnectorInfo.getConnect(),
  106. fullConnectorInfo.getStatus().getState().getValue(),
  107. fullConnectorInfo.getType().getValue());
  108. }
  109. public Mono<ConnectorTopics> getConnectorTopics(KafkaCluster cluster, String connectClusterName,
  110. String connectorName) {
  111. return api(cluster, connectClusterName)
  112. .mono(c -> c.getConnectorTopics(connectorName))
  113. .map(result -> result.get(connectorName))
  114. // old Connect API versions don't have this endpoint, setting empty list for
  115. // backward-compatibility
  116. .onErrorResume(Exception.class, e -> Mono.just(new ConnectorTopics().topics(List.of())));
  117. }
  118. public Flux<String> getConnectorNames(KafkaCluster cluster, String connectName) {
  119. return api(cluster, connectName)
  120. .flux(client -> client.getConnectors(null))
  121. // for some reason `getConnectors` method returns the response as a single string
  122. .collectList().map(e -> e.get(0))
  123. .map(this::parseConnectorsNamesStringToList)
  124. .flatMapMany(Flux::fromIterable);
  125. }
  126. @SneakyThrows
  127. private List<String> parseConnectorsNamesStringToList(String json) {
  128. return objectMapper.readValue(json, new TypeReference<>() {
  129. });
  130. }
  131. public Mono<ConnectorDTO> createConnector(KafkaCluster cluster, String connectName,
  132. Mono<NewConnectorDTO> connector) {
  133. return api(cluster, connectName)
  134. .mono(client ->
  135. connector
  136. .flatMap(c -> connectorExists(cluster, connectName, c.getName())
  137. .map(exists -> {
  138. if (Boolean.TRUE.equals(exists)) {
  139. throw new ValidationException(
  140. String.format("Connector with name %s already exists", c.getName()));
  141. }
  142. return c;
  143. }))
  144. .map(kafkaConnectMapper::toClient)
  145. .flatMap(client::createConnector)
  146. .flatMap(c -> getConnector(cluster, connectName, c.getName()))
  147. );
  148. }
  149. private Mono<Boolean> connectorExists(KafkaCluster cluster, String connectName,
  150. String connectorName) {
  151. return getConnectorNames(cluster, connectName)
  152. .any(name -> name.equals(connectorName));
  153. }
  154. public Mono<ConnectorDTO> getConnector(KafkaCluster cluster, String connectName,
  155. String connectorName) {
  156. return api(cluster, connectName)
  157. .mono(client -> client.getConnector(connectorName)
  158. .map(kafkaConnectMapper::fromClient)
  159. .flatMap(connector ->
  160. client.getConnectorStatus(connector.getName())
  161. // status request can return 404 if tasks not assigned yet
  162. .onErrorResume(WebClientResponseException.NotFound.class,
  163. e -> emptyStatus(connectorName))
  164. .map(connectorStatus -> {
  165. var status = connectorStatus.getConnector();
  166. var sanitizedConfig = kafkaConfigSanitizer.sanitizeConnectorConfig(connector.getConfig());
  167. ConnectorDTO result = new ConnectorDTO()
  168. .connect(connectName)
  169. .status(kafkaConnectMapper.fromClient(status))
  170. .type(connector.getType())
  171. .tasks(connector.getTasks())
  172. .name(connector.getName())
  173. .config(sanitizedConfig);
  174. if (connectorStatus.getTasks() != null) {
  175. boolean isAnyTaskFailed = connectorStatus.getTasks().stream()
  176. .map(TaskStatus::getState)
  177. .anyMatch(TaskStatus.StateEnum.FAILED::equals);
  178. if (isAnyTaskFailed) {
  179. result.getStatus().state(ConnectorStateDTO.TASK_FAILED);
  180. }
  181. }
  182. return result;
  183. })
  184. )
  185. );
  186. }
  187. private Mono<ConnectorStatus> emptyStatus(String connectorName) {
  188. return Mono.just(new ConnectorStatus()
  189. .name(connectorName)
  190. .tasks(List.of())
  191. .connector(new ConnectorStatusConnector()
  192. .state(ConnectorStatusConnector.StateEnum.UNASSIGNED)));
  193. }
  194. public Mono<Map<String, Object>> getConnectorConfig(KafkaCluster cluster, String connectName,
  195. String connectorName) {
  196. return api(cluster, connectName)
  197. .mono(c -> c.getConnectorConfig(connectorName))
  198. .map(kafkaConfigSanitizer::sanitizeConnectorConfig);
  199. }
  200. public Mono<ConnectorDTO> setConnectorConfig(KafkaCluster cluster, String connectName,
  201. String connectorName, Mono<Map<String, Object>> requestBody) {
  202. return api(cluster, connectName)
  203. .mono(c ->
  204. requestBody
  205. .flatMap(body -> c.setConnectorConfig(connectorName, body))
  206. .map(kafkaConnectMapper::fromClient));
  207. }
  208. public Mono<Void> deleteConnector(
  209. KafkaCluster cluster, String connectName, String connectorName) {
  210. return api(cluster, connectName)
  211. .mono(c -> c.deleteConnector(connectorName));
  212. }
  213. public Mono<Void> updateConnectorState(KafkaCluster cluster, String connectName,
  214. String connectorName, ConnectorActionDTO action) {
  215. return api(cluster, connectName)
  216. .mono(client -> {
  217. switch (action) {
  218. case RESTART:
  219. return client.restartConnector(connectorName, false, false);
  220. case RESTART_ALL_TASKS:
  221. return restartTasks(cluster, connectName, connectorName, task -> true);
  222. case RESTART_FAILED_TASKS:
  223. return restartTasks(cluster, connectName, connectorName,
  224. t -> t.getStatus().getState() == ConnectorTaskStatusDTO.FAILED);
  225. case PAUSE:
  226. return client.pauseConnector(connectorName);
  227. case RESUME:
  228. return client.resumeConnector(connectorName);
  229. default:
  230. throw new IllegalStateException("Unexpected value: " + action);
  231. }
  232. });
  233. }
  234. private Mono<Void> restartTasks(KafkaCluster cluster, String connectName,
  235. String connectorName, Predicate<TaskDTO> taskFilter) {
  236. return getConnectorTasks(cluster, connectName, connectorName)
  237. .filter(taskFilter)
  238. .flatMap(t ->
  239. restartConnectorTask(cluster, connectName, connectorName, t.getId().getTask()))
  240. .then();
  241. }
  242. public Flux<TaskDTO> getConnectorTasks(KafkaCluster cluster, String connectName, String connectorName) {
  243. return api(cluster, connectName)
  244. .flux(client ->
  245. client.getConnectorTasks(connectorName)
  246. .onErrorResume(WebClientResponseException.NotFound.class, e -> Flux.empty())
  247. .map(kafkaConnectMapper::fromClient)
  248. .flatMap(task ->
  249. client
  250. .getConnectorTaskStatus(connectorName, task.getId().getTask())
  251. .onErrorResume(WebClientResponseException.NotFound.class, e -> Mono.empty())
  252. .map(kafkaConnectMapper::fromClient)
  253. .map(task::status)
  254. ));
  255. }
  256. public Mono<Void> restartConnectorTask(KafkaCluster cluster, String connectName,
  257. String connectorName, Integer taskId) {
  258. return api(cluster, connectName)
  259. .mono(client -> client.restartConnectorTask(connectorName, taskId));
  260. }
  261. public Flux<ConnectorPluginDTO> getConnectorPlugins(KafkaCluster cluster,
  262. String connectName) {
  263. return api(cluster, connectName)
  264. .flux(client -> client.getConnectorPlugins().map(kafkaConnectMapper::fromClient));
  265. }
  266. public Mono<ConnectorPluginConfigValidationResponseDTO> validateConnectorPluginConfig(
  267. KafkaCluster cluster, String connectName, String pluginName, Mono<Map<String, Object>> requestBody) {
  268. return api(cluster, connectName)
  269. .mono(client ->
  270. requestBody
  271. .flatMap(body ->
  272. client.validateConnectorPluginConfig(pluginName, body))
  273. .map(kafkaConnectMapper::fromClient)
  274. );
  275. }
  276. private ReactiveFailover<KafkaConnectClientApi> api(KafkaCluster cluster, String connectName) {
  277. var client = cluster.getConnectsClients().get(connectName);
  278. if (client == null) {
  279. throw new NotFoundException(
  280. "Connect %s not found for cluster %s".formatted(connectName, cluster.getName()));
  281. }
  282. return client;
  283. }
  284. }