KafkaConnectService.java 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319
  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.HashMap;
  27. import java.util.List;
  28. import java.util.Map;
  29. import java.util.Optional;
  30. import java.util.function.Predicate;
  31. import java.util.stream.Collectors;
  32. import java.util.stream.Stream;
  33. import javax.annotation.Nullable;
  34. import lombok.RequiredArgsConstructor;
  35. import lombok.SneakyThrows;
  36. import lombok.extern.slf4j.Slf4j;
  37. import org.apache.commons.lang3.StringUtils;
  38. import org.springframework.stereotype.Service;
  39. import org.springframework.web.reactive.function.client.WebClientResponseException;
  40. import reactor.core.publisher.Flux;
  41. import reactor.core.publisher.Mono;
  42. import reactor.util.function.Tuples;
  43. @Service
  44. @Slf4j
  45. @RequiredArgsConstructor
  46. public class KafkaConnectService {
  47. private final ClusterMapper clusterMapper;
  48. private final KafkaConnectMapper kafkaConnectMapper;
  49. private final ObjectMapper objectMapper;
  50. private final KafkaConfigSanitizer kafkaConfigSanitizer;
  51. public Flux<ConnectDTO> getConnects(KafkaCluster cluster) {
  52. return Flux.fromIterable(
  53. Optional.ofNullable(cluster.getOriginalProperties().getKafkaConnect())
  54. .map(lst -> lst.stream().map(clusterMapper::toKafkaConnect).toList())
  55. .orElse(List.of())
  56. );
  57. }
  58. public Flux<FullConnectorInfoDTO> getAllConnectors(final KafkaCluster cluster,
  59. @Nullable final String search) {
  60. return getConnects(cluster)
  61. .flatMap(connect -> getConnectorNames(cluster, connect.getName()).map(cn -> Tuples.of(connect.getName(), cn)))
  62. .flatMap(pair -> getConnector(cluster, pair.getT1(), pair.getT2()))
  63. .flatMap(connector ->
  64. getConnectorConfig(cluster, connector.getConnect(), connector.getName())
  65. .map(config -> InternalConnectInfo.builder()
  66. .connector(connector)
  67. .config(config)
  68. .build()
  69. )
  70. )
  71. .flatMap(connectInfo -> {
  72. ConnectorDTO connector = connectInfo.getConnector();
  73. return getConnectorTasks(cluster, connector.getConnect(), connector.getName())
  74. .collectList()
  75. .map(tasks -> InternalConnectInfo.builder()
  76. .connector(connector)
  77. .config(connectInfo.getConfig())
  78. .tasks(tasks)
  79. .build()
  80. );
  81. })
  82. .flatMap(connectInfo -> {
  83. ConnectorDTO connector = connectInfo.getConnector();
  84. return getConnectorTopics(cluster, connector.getConnect(), connector.getName())
  85. .map(ct -> InternalConnectInfo.builder()
  86. .connector(connector)
  87. .config(connectInfo.getConfig())
  88. .tasks(connectInfo.getTasks())
  89. .topics(ct.getTopics())
  90. .build()
  91. );
  92. })
  93. .map(kafkaConnectMapper::fullConnectorInfoFromTuple)
  94. .filter(matchesSearchTerm(search));
  95. }
  96. private Predicate<FullConnectorInfoDTO> matchesSearchTerm(@Nullable final String search) {
  97. if (search == null) {
  98. return c -> true;
  99. }
  100. return connector -> getStringsForSearch(connector)
  101. .anyMatch(string -> StringUtils.containsIgnoreCase(string, search));
  102. }
  103. private Stream<String> getStringsForSearch(FullConnectorInfoDTO fullConnectorInfo) {
  104. return Stream.of(
  105. fullConnectorInfo.getName(),
  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. final Map<String, Object> obfuscatedConfig = connector.getConfig().entrySet()
  167. .stream()
  168. .collect(Collectors.toMap(
  169. Map.Entry::getKey,
  170. e -> kafkaConfigSanitizer.sanitize(e.getKey(), e.getValue())
  171. ));
  172. ConnectorDTO result = (ConnectorDTO) new ConnectorDTO()
  173. .connect(connectName)
  174. .status(kafkaConnectMapper.fromClient(status))
  175. .type(connector.getType())
  176. .tasks(connector.getTasks())
  177. .name(connector.getName())
  178. .config(obfuscatedConfig);
  179. if (connectorStatus.getTasks() != null) {
  180. boolean isAnyTaskFailed = connectorStatus.getTasks().stream()
  181. .map(TaskStatus::getState)
  182. .anyMatch(TaskStatus.StateEnum.FAILED::equals);
  183. if (isAnyTaskFailed) {
  184. result.getStatus().state(ConnectorStateDTO.TASK_FAILED);
  185. }
  186. }
  187. return result;
  188. })
  189. )
  190. );
  191. }
  192. private Mono<ConnectorStatus> emptyStatus(String connectorName) {
  193. return Mono.just(new ConnectorStatus()
  194. .name(connectorName)
  195. .tasks(List.of())
  196. .connector(new ConnectorStatusConnector()
  197. .state(ConnectorStatusConnector.StateEnum.UNASSIGNED)));
  198. }
  199. public Mono<Map<String, Object>> getConnectorConfig(KafkaCluster cluster, String connectName,
  200. String connectorName) {
  201. return api(cluster, connectName)
  202. .mono(c -> c.getConnectorConfig(connectorName))
  203. .map(connectorConfig -> {
  204. final Map<String, Object> obfuscatedMap = new HashMap<>();
  205. connectorConfig.forEach((key, value) ->
  206. obfuscatedMap.put(key, kafkaConfigSanitizer.sanitize(key, value)));
  207. return obfuscatedMap;
  208. });
  209. }
  210. public Mono<ConnectorDTO> setConnectorConfig(KafkaCluster cluster, String connectName,
  211. String connectorName, Mono<Map<String, Object>> requestBody) {
  212. return api(cluster, connectName)
  213. .mono(c ->
  214. requestBody
  215. .flatMap(body -> c.setConnectorConfig(connectorName, body))
  216. .map(kafkaConnectMapper::fromClient));
  217. }
  218. public Mono<Void> deleteConnector(
  219. KafkaCluster cluster, String connectName, String connectorName) {
  220. return api(cluster, connectName)
  221. .mono(c -> c.deleteConnector(connectorName));
  222. }
  223. public Mono<Void> updateConnectorState(KafkaCluster cluster, String connectName,
  224. String connectorName, ConnectorActionDTO action) {
  225. return api(cluster, connectName)
  226. .mono(client -> {
  227. switch (action) {
  228. case RESTART:
  229. return client.restartConnector(connectorName, false, false);
  230. case RESTART_ALL_TASKS:
  231. return restartTasks(cluster, connectName, connectorName, task -> true);
  232. case RESTART_FAILED_TASKS:
  233. return restartTasks(cluster, connectName, connectorName,
  234. t -> t.getStatus().getState() == ConnectorTaskStatusDTO.FAILED);
  235. case PAUSE:
  236. return client.pauseConnector(connectorName);
  237. case RESUME:
  238. return client.resumeConnector(connectorName);
  239. default:
  240. throw new IllegalStateException("Unexpected value: " + action);
  241. }
  242. });
  243. }
  244. private Mono<Void> restartTasks(KafkaCluster cluster, String connectName,
  245. String connectorName, Predicate<TaskDTO> taskFilter) {
  246. return getConnectorTasks(cluster, connectName, connectorName)
  247. .filter(taskFilter)
  248. .flatMap(t ->
  249. restartConnectorTask(cluster, connectName, connectorName, t.getId().getTask()))
  250. .then();
  251. }
  252. public Flux<TaskDTO> getConnectorTasks(KafkaCluster cluster, String connectName, String connectorName) {
  253. return api(cluster, connectName)
  254. .flux(client ->
  255. client.getConnectorTasks(connectorName)
  256. .onErrorResume(WebClientResponseException.NotFound.class, e -> Flux.empty())
  257. .map(kafkaConnectMapper::fromClient)
  258. .flatMap(task ->
  259. client
  260. .getConnectorTaskStatus(connectorName, task.getId().getTask())
  261. .onErrorResume(WebClientResponseException.NotFound.class, e -> Mono.empty())
  262. .map(kafkaConnectMapper::fromClient)
  263. .map(task::status)
  264. ));
  265. }
  266. public Mono<Void> restartConnectorTask(KafkaCluster cluster, String connectName,
  267. String connectorName, Integer taskId) {
  268. return api(cluster, connectName)
  269. .mono(client -> client.restartConnectorTask(connectorName, taskId));
  270. }
  271. public Flux<ConnectorPluginDTO> getConnectorPlugins(KafkaCluster cluster,
  272. String connectName) {
  273. return api(cluster, connectName)
  274. .flux(client -> client.getConnectorPlugins().map(kafkaConnectMapper::fromClient));
  275. }
  276. public Mono<ConnectorPluginConfigValidationResponseDTO> validateConnectorPluginConfig(
  277. KafkaCluster cluster, String connectName, String pluginName, Mono<Map<String, Object>> requestBody) {
  278. return api(cluster, connectName)
  279. .mono(client ->
  280. requestBody
  281. .flatMap(body ->
  282. client.validateConnectorPluginConfig(pluginName, body))
  283. .map(kafkaConnectMapper::fromClient)
  284. );
  285. }
  286. private ReactiveFailover<KafkaConnectClientApi> api(KafkaCluster cluster, String connectName) {
  287. var client = cluster.getConnectsClients().get(connectName);
  288. if (client == null) {
  289. throw new NotFoundException(
  290. "Connect %s not found for cluster %s".formatted(connectName, cluster.getName()));
  291. }
  292. return client;
  293. }
  294. }