KafkaConnectService.java 12 KB

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