KafkaConnectService.java 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307
  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.client.KafkaConnectClients;
  5. import com.provectus.kafka.ui.connect.model.ConnectorTopics;
  6. import com.provectus.kafka.ui.connect.model.TaskStatus;
  7. import com.provectus.kafka.ui.exception.ClusterNotFoundException;
  8. import com.provectus.kafka.ui.exception.ConnectNotFoundException;
  9. import com.provectus.kafka.ui.mapper.ClusterMapper;
  10. import com.provectus.kafka.ui.mapper.KafkaConnectMapper;
  11. import com.provectus.kafka.ui.model.Connect;
  12. import com.provectus.kafka.ui.model.Connector;
  13. import com.provectus.kafka.ui.model.ConnectorAction;
  14. import com.provectus.kafka.ui.model.ConnectorPlugin;
  15. import com.provectus.kafka.ui.model.ConnectorPluginConfigValidationResponse;
  16. import com.provectus.kafka.ui.model.ConnectorState;
  17. import com.provectus.kafka.ui.model.FullConnectorInfo;
  18. import com.provectus.kafka.ui.model.KafkaCluster;
  19. import com.provectus.kafka.ui.model.KafkaConnectCluster;
  20. import com.provectus.kafka.ui.model.NewConnector;
  21. import com.provectus.kafka.ui.model.Task;
  22. import com.provectus.kafka.ui.model.connect.InternalConnectInfo;
  23. import java.util.Collection;
  24. import java.util.List;
  25. import java.util.Map;
  26. import java.util.function.Function;
  27. import java.util.function.Predicate;
  28. import java.util.stream.Collectors;
  29. import java.util.stream.Stream;
  30. import lombok.RequiredArgsConstructor;
  31. import lombok.SneakyThrows;
  32. import lombok.extern.log4j.Log4j2;
  33. import org.apache.commons.lang3.StringUtils;
  34. import org.springframework.stereotype.Service;
  35. import reactor.core.publisher.Flux;
  36. import reactor.core.publisher.Mono;
  37. import reactor.util.function.Tuple2;
  38. import reactor.util.function.Tuples;
  39. @Service
  40. @Log4j2
  41. @RequiredArgsConstructor
  42. public class KafkaConnectService {
  43. private final ClustersStorage clustersStorage;
  44. private final ClusterMapper clusterMapper;
  45. private final KafkaConnectMapper kafkaConnectMapper;
  46. private final ObjectMapper objectMapper;
  47. public Mono<Flux<Connect>> getConnects(String clusterName) {
  48. return Mono.just(
  49. Flux.fromIterable(clustersStorage.getClusterByName(clusterName)
  50. .map(KafkaCluster::getKafkaConnect).stream()
  51. .flatMap(Collection::stream)
  52. .map(clusterMapper::toKafkaConnect)
  53. .collect(Collectors.toList())
  54. )
  55. );
  56. }
  57. public Flux<FullConnectorInfo> getAllConnectors(final String clusterName, final String search) {
  58. return getConnects(clusterName)
  59. .flatMapMany(Function.identity())
  60. .flatMap(connect -> getConnectorNames(clusterName, connect))
  61. .flatMap(pair -> getConnector(clusterName, pair.getT1(), pair.getT2()))
  62. .flatMap(connector ->
  63. getConnectorConfig(clusterName, connector.getConnect(), connector.getName())
  64. .map(config -> InternalConnectInfo.builder()
  65. .connector(connector)
  66. .config(config)
  67. .build()
  68. )
  69. )
  70. .flatMap(connectInfo -> {
  71. Connector connector = connectInfo.getConnector();
  72. return getConnectorTasks(clusterName, 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. Connector connector = connectInfo.getConnector();
  83. return getConnectorTopics(clusterName, 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<FullConnectorInfo> matchesSearchTerm(final String search) {
  96. return (connector) -> getSearchValues(connector)
  97. .anyMatch(value -> value.contains(
  98. StringUtils.defaultString(
  99. search,
  100. StringUtils.EMPTY)
  101. .toUpperCase()));
  102. }
  103. private Stream<String> getSearchValues(FullConnectorInfo fullConnectorInfo) {
  104. return Stream.of(
  105. fullConnectorInfo.getName(),
  106. fullConnectorInfo.getStatus().getState().getValue(),
  107. fullConnectorInfo.getType().getValue())
  108. .map(String::toUpperCase);
  109. }
  110. private Mono<ConnectorTopics> getConnectorTopics(String clusterName, String connectClusterName,
  111. String connectorName) {
  112. return getConnectAddress(clusterName, connectClusterName)
  113. .flatMap(connectUrl -> KafkaConnectClients
  114. .withBaseUrl(connectUrl)
  115. .getConnectorTopics(connectorName)
  116. .map(result -> result.get(connectorName))
  117. );
  118. }
  119. private Flux<Tuple2<String, String>> getConnectorNames(String clusterName, Connect connect) {
  120. return getConnectors(clusterName, connect.getName())
  121. .collectList().map(e -> e.get(0))
  122. // for some reason `getConnectors` method returns the response as a single string
  123. .map(this::parseToList)
  124. .flatMapMany(Flux::fromIterable)
  125. .map(connector -> Tuples.of(connect.getName(), connector));
  126. }
  127. @SneakyThrows
  128. private List<String> parseToList(String json) {
  129. return objectMapper.readValue(json, new TypeReference<>() {
  130. });
  131. }
  132. public Flux<String> getConnectors(String clusterName, String connectName) {
  133. return getConnectAddress(clusterName, connectName)
  134. .flatMapMany(connect ->
  135. KafkaConnectClients.withBaseUrl(connect).getConnectors(null)
  136. .doOnError(log::error)
  137. );
  138. }
  139. public Mono<Connector> createConnector(String clusterName, String connectName,
  140. Mono<NewConnector> connector) {
  141. return getConnectAddress(clusterName, connectName)
  142. .flatMap(connect ->
  143. connector
  144. .map(kafkaConnectMapper::toClient)
  145. .flatMap(c ->
  146. KafkaConnectClients.withBaseUrl(connect).createConnector(c)
  147. )
  148. .flatMap(c -> getConnector(clusterName, connectName, c.getName()))
  149. );
  150. }
  151. public Mono<Connector> getConnector(String clusterName, String connectName,
  152. String connectorName) {
  153. return getConnectAddress(clusterName, connectName)
  154. .flatMap(connect -> KafkaConnectClients.withBaseUrl(connect).getConnector(connectorName)
  155. .map(kafkaConnectMapper::fromClient)
  156. .flatMap(connector ->
  157. KafkaConnectClients.withBaseUrl(connect).getConnectorStatus(connector.getName())
  158. .map(connectorStatus -> {
  159. var status = connectorStatus.getConnector();
  160. Connector result = (Connector) new Connector()
  161. .connect(connectName)
  162. .status(kafkaConnectMapper.fromClient(status))
  163. .type(connector.getType())
  164. .tasks(connector.getTasks())
  165. .name(connector.getName())
  166. .config(connector.getConfig());
  167. if (connectorStatus.getTasks() != null) {
  168. boolean isAnyTaskFailed = connectorStatus.getTasks().stream()
  169. .map(TaskStatus::getState)
  170. .anyMatch(TaskStatus.StateEnum.FAILED::equals);
  171. if (isAnyTaskFailed) {
  172. result.getStatus().state(ConnectorState.TASK_FAILED);
  173. }
  174. }
  175. return result;
  176. })
  177. )
  178. );
  179. }
  180. public Mono<Map<String, Object>> getConnectorConfig(String clusterName, String connectName,
  181. String connectorName) {
  182. return getConnectAddress(clusterName, connectName)
  183. .flatMap(connect ->
  184. KafkaConnectClients.withBaseUrl(connect).getConnectorConfig(connectorName)
  185. );
  186. }
  187. public Mono<Connector> setConnectorConfig(String clusterName, String connectName,
  188. String connectorName, Mono<Object> requestBody) {
  189. return getConnectAddress(clusterName, connectName)
  190. .flatMap(connect ->
  191. requestBody.flatMap(body ->
  192. KafkaConnectClients.withBaseUrl(connect)
  193. .setConnectorConfig(connectorName, (Map<String, Object>) body)
  194. )
  195. .map(kafkaConnectMapper::fromClient)
  196. );
  197. }
  198. public Mono<Void> deleteConnector(String clusterName, String connectName, String connectorName) {
  199. return getConnectAddress(clusterName, connectName)
  200. .flatMap(connect ->
  201. KafkaConnectClients.withBaseUrl(connect).deleteConnector(connectorName)
  202. );
  203. }
  204. public Mono<Void> updateConnectorState(String clusterName, String connectName,
  205. String connectorName, ConnectorAction action) {
  206. Function<String, Mono<Void>> kafkaClientCall;
  207. switch (action) {
  208. case RESTART:
  209. kafkaClientCall =
  210. connect -> KafkaConnectClients.withBaseUrl(connect).restartConnector(connectorName);
  211. break;
  212. case PAUSE:
  213. kafkaClientCall =
  214. connect -> KafkaConnectClients.withBaseUrl(connect).pauseConnector(connectorName);
  215. break;
  216. case RESUME:
  217. kafkaClientCall =
  218. connect -> KafkaConnectClients.withBaseUrl(connect).resumeConnector(connectorName);
  219. break;
  220. default:
  221. throw new IllegalStateException("Unexpected value: " + action);
  222. }
  223. return getConnectAddress(clusterName, connectName)
  224. .flatMap(kafkaClientCall);
  225. }
  226. public Flux<Task> getConnectorTasks(String clusterName, String connectName,
  227. String connectorName) {
  228. return getConnectAddress(clusterName, connectName)
  229. .flatMapMany(connect ->
  230. KafkaConnectClients.withBaseUrl(connect).getConnectorTasks(connectorName)
  231. .map(kafkaConnectMapper::fromClient)
  232. .flatMap(task ->
  233. KafkaConnectClients.withBaseUrl(connect)
  234. .getConnectorTaskStatus(connectorName, task.getId().getTask())
  235. .map(kafkaConnectMapper::fromClient)
  236. .map(task::status)
  237. )
  238. );
  239. }
  240. public Mono<Void> restartConnectorTask(String clusterName, String connectName,
  241. String connectorName, Integer taskId) {
  242. return getConnectAddress(clusterName, connectName)
  243. .flatMap(connect ->
  244. KafkaConnectClients.withBaseUrl(connect).restartConnectorTask(connectorName, taskId)
  245. );
  246. }
  247. public Mono<Flux<ConnectorPlugin>> getConnectorPlugins(String clusterName, String connectName) {
  248. return Mono.just(getConnectAddress(clusterName, connectName)
  249. .flatMapMany(connect ->
  250. KafkaConnectClients.withBaseUrl(connect).getConnectorPlugins()
  251. .map(kafkaConnectMapper::fromClient)
  252. ));
  253. }
  254. public Mono<ConnectorPluginConfigValidationResponse> validateConnectorPluginConfig(
  255. String clusterName, String connectName, String pluginName, Mono<Object> requestBody) {
  256. return getConnectAddress(clusterName, connectName)
  257. .flatMap(connect ->
  258. requestBody.flatMap(body ->
  259. KafkaConnectClients.withBaseUrl(connect)
  260. .validateConnectorPluginConfig(pluginName, (Map<String, Object>) body)
  261. )
  262. .map(kafkaConnectMapper::fromClient)
  263. );
  264. }
  265. private Mono<KafkaCluster> getCluster(String clusterName) {
  266. return clustersStorage.getClusterByName(clusterName)
  267. .map(Mono::just)
  268. .orElse(Mono.error(ClusterNotFoundException::new));
  269. }
  270. private Mono<String> getConnectAddress(String clusterName, String connectName) {
  271. return getCluster(clusterName)
  272. .map(kafkaCluster ->
  273. kafkaCluster.getKafkaConnect().stream()
  274. .filter(connect -> connect.getName().equals(connectName))
  275. .findFirst()
  276. .map(KafkaConnectCluster::getAddress)
  277. )
  278. .flatMap(connect -> connect
  279. .map(Mono::just)
  280. .orElse(Mono.error(ConnectNotFoundException::new))
  281. );
  282. }
  283. }