{"id":354025,"date":"2024-05-20T22:40:07","date_gmt":"2024-05-20T22:40:07","guid":{"rendered":"http:\/\/savepearlharbor.com\/?p=354025"},"modified":"-0001-11-30T00:00:00","modified_gmt":"-0001-11-29T21:00:00","slug":"","status":"publish","type":"post","link":"https:\/\/savepearlharbor.com\/?p=354025","title":{"rendered":"<span>\u0414\u0438\u043d\u0430\u043c\u0438\u0447\u0435\u0441\u043a\u043e\u0435 \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u0435 \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u0435\u0439 \u0432 Kafka<\/span>"},"content":{"rendered":"<div><!--[--><!--]--><\/div>\n<div id=\"post-content-body\">\n<div>\n<div class=\"article-formatted-body article-formatted-body article-formatted-body_version-2\">\n<div xmlns=\"http:\/\/www.w3.org\/1999\/xhtml\">\n<p>\u042d\u0442\u0430 \u0441\u0442\u0430\u0442\u044c\u044f \u043e\u0431\u044c\u044f\u0441\u043d\u044f\u0435\u0442, \u043a\u0430\u043a \u0441\u043e\u0437\u0434\u0430\u0442\u044c \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u044c \u0432 Kafka \u043d\u0430 \u043b\u0435\u0442\u0443 \u0432 \u043f\u0440\u043e\u0446\u0435\u0441\u0441\u0435 \u0440\u0430\u0431\u043e\u0442\u044b \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u044f.<\/p>\n<figure class=\"full-width\"><img loading=\"lazy\" decoding=\"async\" src=\"https:\/\/habrastorage.org\/r\/w1560\/getpro\/habr\/upload_files\/58d\/827\/c6d\/58d827c6d71b24f6389f0b54d4fb6fc2.png\" width=\"1200\" height=\"630\" data-src=\"https:\/\/habrastorage.org\/getpro\/habr\/upload_files\/58d\/827\/c6d\/58d827c6d71b24f6389f0b54d4fb6fc2.png\"\/><\/figure>\n<h3>Plan:<\/h3>\n<ol>\n<li>\n<p>\u0421\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u0448\u0430\u0431\u043b\u043e\u043d\u043d\u044b\u0439 \u043a\u043b\u0430\u0441\u0441 \u0447\u0435\u0440\u0435\u0437 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u044e \u0438\u043d\u0442\u0435\u0440\u0444\u0435\u0439\u0441\u0430 MessageListener.<\/p>\n<\/li>\n<li>\n<p>\u0421\u043e\u0437\u0434\u0430\u0434\u0438\u043c KafkaListenerEndpoint \u0441 \u043f\u043e\u043c\u043e\u0449\u044c\u044e \u0448\u0430\u0431\u043b\u043e\u043d\u0430.<\/p>\n<\/li>\n<li>\n<p>\u0417\u0430\u0440\u0435\u0433\u0435\u0441\u0442\u0440\u0438\u0440\u0443\u0435\u043c \u044d\u043d\u0434\u043f\u043e\u0438\u043d\u0442 \u0432 KafkaListenerEndpointRegistry.<\/p>\n<\/li>\n<li>\n<p>\u0421\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u043e\u043a\u0440\u0443\u0436\u0435\u043d\u0438\u0435 \u0434\u043b\u044f \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f.<\/p>\n<\/li>\n<li>\n<p>\u041f\u0440\u043e\u0442\u0435\u0441\u0442\u0438\u0440\u0443\u0435\u043c \u0440\u0435\u0448\u0435\u043d\u0438\u0435.<\/p>\n<\/li>\n<li>\n<p>\u0417\u0430\u043a\u043b\u044e\u0447\u0435\u043d\u0438\u0435.<\/p>\n<\/li>\n<\/ol>\n<h3>1. \u0421\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u0448\u0430\u0431\u043b\u043e\u043d\u043d\u044b\u0439 \u043a\u043b\u0430\u0441\u0441 \u0447\u0435\u0440\u0435\u0437 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u044e \u0438\u043d\u0442\u0435\u0440\u0444\u0435\u0439\u0441\u0430 MessageListener<\/h3>\n<p>\u0421\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u043a\u043b\u0430\u0441\u0441\u00a0<em>KafkaTemplateListener<\/em>\u00a0\u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u0440\u0435\u0430\u043b\u0438\u0437\u0443\u0435\u0442 \u0438\u043d\u0442\u0435\u0440\u0444\u0435\u0439\u0441\u00a0<em>MessageListener<\/em>. \u042d\u0442\u043e\u0442 \u0448\u0430\u0431\u043b\u043e\u043d \u0438\u0441\u0442\u043e\u0447\u043d\u0438\u043a \u043b\u043e\u0433\u0438\u043a\u0438 \u0434\u043b\u044f \u0431\u0443\u0434\u0443\u0449\u0438\u0445 \u0434\u0438\u043d\u0430\u043c\u0438\u0447\u0435\u0441\u043a\u0438 \u0441\u043e\u0437\u0434\u0430\u043d\u043d\u044b\u0445 \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u0435\u0439.<\/p>\n<pre><code class=\"java\">public class KafkaTemplateListener implements MessageListener&lt;String, String> {      @Override         public void onMessage(ConsumerRecord&lt;String, String> record) {         System.out.println(\"RECORD PROCESSING: \" + record);     } }<\/code><\/pre>\n<h3>2. \u0421\u043e\u0437\u0434\u0430\u0434\u0438\u043c KafkaListenerEndpoint \u0441 \u043f\u043e\u043c\u043e\u0449\u044c\u044e \u0440\u0435\u0430\u043b\u0438\u0437\u043e\u0432\u0430\u043d\u043d\u043e\u0433\u043e \u0448\u0430\u0431\u043b\u043e\u043d\u0430<\/h3>\n<p>\u0412 \u043c\u0435\u0442\u043e\u0434\u0435\u00a0<em>createDefaultMethodKafkaListenerEndpoint(String topic)\u00a0<\/em>\u043d\u0443\u0436\u043d\u043e \u0443\u0441\u0442\u0430\u043d\u043e\u0432\u0438\u0442\u044c \u043d\u0430\u0441\u0442\u0440\u043e\u0439\u043a\u0438, \u0442\u0430\u043a\u0438\u0435 \u043a\u0430\u043a Endpoint Id, Group Id, Topics \u0438 \u0442.\u0434.<\/p>\n<p>\u0412 \u043c\u0435\u0442\u043e\u0434\u0435\u00a0<em>createKafkaListenerEndpoint(String topic)\u00a0<\/em>\u043d\u0443\u0436\u043d\u043e \u0443\u0441\u0442\u0430\u043d\u043e\u0432\u0438\u0442\u044c \u0448\u0430\u0431\u043b\u043e\u043d \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u044f \u0438 \u043c\u0435\u0442\u043e\u0434, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u0441\u043b\u0443\u0448\u0430\u0435\u0442 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0435 \u0438\u0437 Kafka<\/p>\n<pre><code class=\"java\">@Service public class KafkaListenerCreator {     String kafkaGroupId = \"kafkaGroupId\";     String kafkaListenerId = \"kafkaListenerId-\";     static AtomicLong endpointIdIndex = new AtomicLong(1);      private KafkaListenerEndpoint createKafkaListenerEndpoint(String topic) {         MethodKafkaListenerEndpoint&lt;String, String> kafkaListenerEndpoint =             createDefaultMethodKafkaListenerEndpoint(topic);         kafkaListenerEndpoint.setBean(new KafkaTemplateListener());         try {             kafkaListenerEndpoint.setMethod(KafkaTemplateListener.class.getMethod(\"onMessage\", ConsumerRecord.class));         } catch (NoSuchMethodException e) {             throw new RuntimeException(\"Attempt to call a non-existent method \" + e);         }         return kafkaListenerEndpoint;     }      private MethodKafkaListenerEndpoint&lt;String, String> createDefaultMethodKafkaListenerEndpoint(String topic) {         MethodKafkaListenerEndpoint&lt;String, String> kafkaListenerEndpoint = new MethodKafkaListenerEndpoint&lt;>();         kafkaListenerEndpoint.setId(generateListenerId());         kafkaListenerEndpoint.setGroupId(kafkaGroupId);         kafkaListenerEndpoint.setAutoStartup(true);         kafkaListenerEndpoint.setTopics(topic);         kafkaListenerEndpoint.setMessageHandlerMethodFactory(new DefaultMessageHandlerMethodFactory());         return kafkaListenerEndpoint;     }      private String generateListenerId() {         return kafkaGeneralListenerEndpointId + endpointIdIndex.getAndIncrement();      } }<\/code><\/pre>\n<h3>3. \u0417\u0430\u0440\u0435\u0433\u0435\u0441\u0442\u0440\u0438\u0440\u0443\u0435\u043c \u044d\u043d\u0434\u043f\u043e\u0438\u043d\u0442 \u0432 KafkaListenerEndpointRegistry<\/h3>\n<p>\u041b\u043e\u0433\u0438\u043a\u0430, \u043a\u043e\u0442\u043e\u0440\u0430\u044f \u0440\u0435\u0433\u0438\u0441\u0442\u0440\u0438\u0440\u0443\u0435\u0442 \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u044c \u043d\u0430\u0445\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u0432 \u0442\u0435\u043b\u0435 \u043c\u0435\u0442\u043e\u0434\u0430 <em>createAndRegisterListener(String topic).\u00a0<\/em>\u041b\u043e\u0433\u0438\u043a\u0430 \u043d\u0430\u0445\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u0432 \u0442\u043e\u043c \u0436\u0435 \u043a\u043b\u0430\u0441\u0441\u0435<em>\u00a0KafkaListenerCreator.<\/em><\/p>\n<pre><code class=\"java\">@Service public class KafkaListenerCreator {   \/\/... HERE HAS TO BE VARIABLES FROM PREVIOUS EXAMPLE    @Autowired   private KafkaListenerEndpointRegistry kafkaListenerEndpointRegistry;   @Autowired   private KafkaListenerContainerFactory kafkaListenerContainerFactory;    public void createAndRegisterListener(String topic) {     KafkaListenerEndpoint listener = createKafkaListenerEndpoint(topic);     kafkaListenerEndpointRegistry.registerListenerContainer(listener, kafkaListenerContainerFactory, true);   }    \/\/... HERE HAS TO BE METHODS FROM PREVIOUS EXAMPLE  }<\/code><\/pre>\n<h3>4. \u0421\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u043e\u043a\u0440\u0443\u0436\u0435\u043d\u0438\u0435 \u0434\u043b\u044f \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f<\/h3>\n<p>\u0421\u043d\u0430\u0447\u0430\u043b\u0430 \u0441\u043e\u0437\u0434\u0430\u0434\u0438\u043c REST \u043c\u0435\u0442\u043e\u0434 \u0434\u043b\u044f \u0434\u043b\u044f \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u044f \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u044f Kafka.<\/p>\n<p>\u042f \u0441\u043e\u0437\u0434\u0430\u043b \u043a\u043b\u0430\u0441\u0441 <em>KafkaController \u0438 \u043c\u0435\u0442\u043e\u0434 create(String topic). <\/em>\u042d\u0442\u043e\u0442 \u043c\u0435\u0442\u043e\u0434 \u043c\u043e\u0436\u0435\u0442 \u0431\u044b\u0442\u044c \u0432\u044b\u0437\u0432\u0430\u043d \u0441 \u043f\u043e\u043c\u043e\u0449\u044c\u044e POST HTTP<em> \u0437\u0430\u043f\u0440\u043e\u0441\u0430.<\/em><\/p>\n<pre><code class=\"java\">@RestController public class KafkaController {     @Autowired     KafkaListenerCreator kafkaListenerCreator;      @PostMapping(path = \"\/create\")     @ResponseStatus(HttpStatus.OK)     public void create(@RequestParam String topic) {         kafkaListenerCreator.createAndRegisterListener(topic);     } }<\/code><\/pre>\n<p>\u0414\u0430\u043b\u0435\u0435 \u0441\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u043c\u0435\u0442\u043e\u0434 \u0434\u043b\u044f \u043e\u0442\u043f\u0440\u0430\u0432\u043a\u0438 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f \u0432 \u043d\u043e\u0432\u043e \u0441\u043e\u0437\u0434\u0430\u043d\u043d\u044b\u0439 \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u044c.<\/p>\n<p>\u041c\u0435\u0442\u043e\u0434\u00a0<em>send(@RequestParam String topic, @RequestParam String message)\u00a0<\/em>\u0438\u043c\u0435\u0435\u0442 \u0434\u0432\u0430 \u043f\u0430\u0440\u0430\u043c\u0435\u0442\u0440\u0430. \u0413\u0434\u0435\u00a0<em>\u201ctopic\u201d<\/em>\u00a0\u044d\u0442\u043e \u0438\u043c\u044f \u0442\u043e\u043f\u0438\u043a\u0430 \u0443 \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u044f Kafka, \u0430\u00a0<em>\u201cmessage\u201d<\/em>\u00a0\u044d\u0442\u043e \u0442\u0435\u043a\u0441\u0442 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f, \u043a\u043e\u0442\u043e\u0440\u043e\u0435 \u043e\u0442\u043f\u0440\u0430\u0432\u043b\u044f\u0435\u0442\u0441\u044f \u0447\u0435\u0440\u0435\u0437 <em>KafkaTemplate<\/em>.<\/p>\n<pre><code class=\"java\">@RestController public class KafkaController {     \/\/... HERE HAS TO BE VARIABLES FROM PREVIOUS EXAMPLE     @Autowired     private KafkaTemplate&lt;String, String> kafkaTemplate;      @PostMapping(path = \"\/send\")     @ResponseStatus(HttpStatus.OK)     public void send(@RequestParam String topic, @RequestParam String message) {         kafkaTemplate.send(topic, message);     }     \/\/... HERE HAS TO BE METHODS FROM PREVIOUS EXAMPLE }<\/code><\/pre>\n<h3>5. \u041f\u0440\u043e\u0442\u0435\u0441\u0442\u0438\u0440\u0443\u0435\u043c \u0440\u0435\u0448\u0435\u043d\u0438\u0435<\/h3>\n<ol>\n<li>\n<p>\u0432\u044b\u0437\u0432\u0430\u0442\u044c http:\/\/localhost:8080\/create?topic=myTopic1<\/p>\n<\/li>\n<li>\n<p>\u0432\u044b\u0437\u0432\u0430\u0442\u044c http:\/\/localhost:8080\/send?topic=myTopic1&amp;message=myTxt1<\/p>\n<\/li>\n<li>\n<p>\u043e\u0436\u0438\u0434\u0430\u0435\u043c\u044b\u0439 \u043b\u043e\u0433: \u201c RECORD PROCESSING: myTxt1\u201d<\/p>\n<\/li>\n<\/ol>\n<h3>6. \u0417\u0430\u043a\u043b\u044e\u0447\u0435\u043d\u0438\u0435<\/h3>\n<p>\u042d\u0442\u0430 \u0441\u0442\u0430\u0442\u044c\u044f \u043f\u0440\u0435\u0434\u043e\u0441\u0442\u0430\u0432\u043b\u044f\u0435\u0442 \u0431\u044b\u0441\u0442\u0440\u043e\u0435 \u0440\u0435\u0448\u0435\u043d\u0438\u0435 \u0434\u043b\u044f \u043f\u0440\u043e\u0431\u043b\u0435\u043c\u044b \u0434\u0438\u043d\u0430\u043c\u0438\u0447\u0435\u0441\u043a\u043e\u0433\u043e \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u044f \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u0435\u0439 \u0432 Kafka.<\/p>\n<h2>\u0414\u043e\u043f\u043e\u043b\u043d\u0438\u0442\u0435\u043b\u044c\u043d\u044b\u0435 \u0440\u0435\u0441\u0443\u0440\u0441\u044b<\/h2>\n<p><a href=\"https:\/\/medium.com\/bliblidotcom-techblog\/dynamic-spring-boot-kafka-consumer-af8740f2c703?source=post_page-----4f8f359d715e--------------------------------\" rel=\"noopener noreferrer nofollow\">https:\/\/medium.com\/bliblidotcom-techblog\/dynamic-spring-boot-kafka-consumer-af8740f2c703?source=post_page&#8212;&#8212;4f8f359d715e&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;<\/a><\/p>\n<p><a href=\"https:\/\/github.com\/spring-projects\/spring-kafka\/issues\/460?source=post_page-----4f8f359d715e--------------------------------\" rel=\"noopener noreferrer nofollow\">https:\/\/github.com\/spring-projects\/spring-kafka\/issues\/460?source=post_page&#8212;&#8212;4f8f359d715e&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;<\/a><\/p>\n<\/p>\n<\/div>\n<\/div>\n<\/div>\n<p><!----><!----><\/div>\n<p><!----><!----><br \/> \u0441\u0441\u044b\u043b\u043a\u0430 \u043d\u0430 \u043e\u0440\u0438\u0433\u0438\u043d\u0430\u043b \u0441\u0442\u0430\u0442\u044c\u0438 <a href=\"https:\/\/habr.com\/ru\/articles\/758958\/\"> https:\/\/habr.com\/ru\/articles\/758958\/<\/a><\/p>\n","protected":false},"excerpt":{"rendered":"<div><!--[--><!--]--><\/div>\n<div id=\"post-content-body\">\n<div>\n<div class=\"article-formatted-body article-formatted-body article-formatted-body_version-2\">\n<div xmlns=\"http:\/\/www.w3.org\/1999\/xhtml\">\n<p>\u042d\u0442\u0430 \u0441\u0442\u0430\u0442\u044c\u044f \u043e\u0431\u044c\u044f\u0441\u043d\u044f\u0435\u0442, \u043a\u0430\u043a \u0441\u043e\u0437\u0434\u0430\u0442\u044c \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u044c \u0432 Kafka \u043d\u0430 \u043b\u0435\u0442\u0443 \u0432 \u043f\u0440\u043e\u0446\u0435\u0441\u0441\u0435 \u0440\u0430\u0431\u043e\u0442\u044b \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u044f.<\/p>\n<figure class=\"full-width\"><\/figure>\n<h3>Plan:<\/h3>\n<ol>\n<li>\n<p>\u0421\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u0448\u0430\u0431\u043b\u043e\u043d\u043d\u044b\u0439 \u043a\u043b\u0430\u0441\u0441 \u0447\u0435\u0440\u0435\u0437 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u044e \u0438\u043d\u0442\u0435\u0440\u0444\u0435\u0439\u0441\u0430 MessageListener.<\/p>\n<\/li>\n<li>\n<p>\u0421\u043e\u0437\u0434\u0430\u0434\u0438\u043c KafkaListenerEndpoint \u0441 \u043f\u043e\u043c\u043e\u0449\u044c\u044e \u0448\u0430\u0431\u043b\u043e\u043d\u0430.<\/p>\n<\/li>\n<li>\n<p>\u0417\u0430\u0440\u0435\u0433\u0435\u0441\u0442\u0440\u0438\u0440\u0443\u0435\u043c \u044d\u043d\u0434\u043f\u043e\u0438\u043d\u0442 \u0432 KafkaListenerEndpointRegistry.<\/p>\n<\/li>\n<li>\n<p>\u0421\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u043e\u043a\u0440\u0443\u0436\u0435\u043d\u0438\u0435 \u0434\u043b\u044f \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f.<\/p>\n<\/li>\n<li>\n<p>\u041f\u0440\u043e\u0442\u0435\u0441\u0442\u0438\u0440\u0443\u0435\u043c \u0440\u0435\u0448\u0435\u043d\u0438\u0435.<\/p>\n<\/li>\n<li>\n<p>\u0417\u0430\u043a\u043b\u044e\u0447\u0435\u043d\u0438\u0435.<\/p>\n<\/li>\n<\/ol>\n<h3>1. \u0421\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u0448\u0430\u0431\u043b\u043e\u043d\u043d\u044b\u0439 \u043a\u043b\u0430\u0441\u0441 \u0447\u0435\u0440\u0435\u0437 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u044e \u0438\u043d\u0442\u0435\u0440\u0444\u0435\u0439\u0441\u0430 MessageListener<\/h3>\n<p>\u0421\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u043a\u043b\u0430\u0441\u0441\u00a0<em>KafkaTemplateListener<\/em>\u00a0\u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u0440\u0435\u0430\u043b\u0438\u0437\u0443\u0435\u0442 \u0438\u043d\u0442\u0435\u0440\u0444\u0435\u0439\u0441\u00a0<em>MessageListener<\/em>. \u042d\u0442\u043e\u0442 \u0448\u0430\u0431\u043b\u043e\u043d \u0438\u0441\u0442\u043e\u0447\u043d\u0438\u043a \u043b\u043e\u0433\u0438\u043a\u0438 \u0434\u043b\u044f \u0431\u0443\u0434\u0443\u0449\u0438\u0445 \u0434\u0438\u043d\u0430\u043c\u0438\u0447\u0435\u0441\u043a\u0438 \u0441\u043e\u0437\u0434\u0430\u043d\u043d\u044b\u0445 \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u0435\u0439.<\/p>\n<pre><code class=\"java\">public class KafkaTemplateListener implements MessageListener&lt;String, String> {      @Override         public void onMessage(ConsumerRecord&lt;String, String> record) {         System.out.println(\"RECORD PROCESSING: \" + record);     } }<\/code><\/pre>\n<h3>2. \u0421\u043e\u0437\u0434\u0430\u0434\u0438\u043c KafkaListenerEndpoint \u0441 \u043f\u043e\u043c\u043e\u0449\u044c\u044e \u0440\u0435\u0430\u043b\u0438\u0437\u043e\u0432\u0430\u043d\u043d\u043e\u0433\u043e \u0448\u0430\u0431\u043b\u043e\u043d\u0430<\/h3>\n<p>\u0412 \u043c\u0435\u0442\u043e\u0434\u0435\u00a0<em>createDefaultMethodKafkaListenerEndpoint(String topic)\u00a0<\/em>\u043d\u0443\u0436\u043d\u043e \u0443\u0441\u0442\u0430\u043d\u043e\u0432\u0438\u0442\u044c \u043d\u0430\u0441\u0442\u0440\u043e\u0439\u043a\u0438, \u0442\u0430\u043a\u0438\u0435 \u043a\u0430\u043a Endpoint Id, Group Id, Topics \u0438 \u0442.\u0434.<\/p>\n<p>\u0412 \u043c\u0435\u0442\u043e\u0434\u0435\u00a0<em>createKafkaListenerEndpoint(String topic)\u00a0<\/em>\u043d\u0443\u0436\u043d\u043e \u0443\u0441\u0442\u0430\u043d\u043e\u0432\u0438\u0442\u044c \u0448\u0430\u0431\u043b\u043e\u043d \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u044f \u0438 \u043c\u0435\u0442\u043e\u0434, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u0441\u043b\u0443\u0448\u0430\u0435\u0442 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0435 \u0438\u0437 Kafka<\/p>\n<pre><code class=\"java\">@Service public class KafkaListenerCreator {     String kafkaGroupId = \"kafkaGroupId\";     String kafkaListenerId = \"kafkaListenerId-\";     static AtomicLong endpointIdIndex = new AtomicLong(1);      private KafkaListenerEndpoint createKafkaListenerEndpoint(String topic) {         MethodKafkaListenerEndpoint&lt;String, String> kafkaListenerEndpoint =             createDefaultMethodKafkaListenerEndpoint(topic);         kafkaListenerEndpoint.setBean(new KafkaTemplateListener());         try {             kafkaListenerEndpoint.setMethod(KafkaTemplateListener.class.getMethod(\"onMessage\", ConsumerRecord.class));         } catch (NoSuchMethodException e) {             throw new RuntimeException(\"Attempt to call a non-existent method \" + e);         }         return kafkaListenerEndpoint;     }      private MethodKafkaListenerEndpoint&lt;String, String> createDefaultMethodKafkaListenerEndpoint(String topic) {         MethodKafkaListenerEndpoint&lt;String, String> kafkaListenerEndpoint = new MethodKafkaListenerEndpoint&lt;>();         kafkaListenerEndpoint.setId(generateListenerId());         kafkaListenerEndpoint.setGroupId(kafkaGroupId);         kafkaListenerEndpoint.setAutoStartup(true);         kafkaListenerEndpoint.setTopics(topic);         kafkaListenerEndpoint.setMessageHandlerMethodFactory(new DefaultMessageHandlerMethodFactory());         return kafkaListenerEndpoint;     }      private String generateListenerId() {         return kafkaGeneralListenerEndpointId + endpointIdIndex.getAndIncrement();      } }<\/code><\/pre>\n<h3>3. \u0417\u0430\u0440\u0435\u0433\u0435\u0441\u0442\u0440\u0438\u0440\u0443\u0435\u043c \u044d\u043d\u0434\u043f\u043e\u0438\u043d\u0442 \u0432 KafkaListenerEndpointRegistry<\/h3>\n<p>\u041b\u043e\u0433\u0438\u043a\u0430, \u043a\u043e\u0442\u043e\u0440\u0430\u044f \u0440\u0435\u0433\u0438\u0441\u0442\u0440\u0438\u0440\u0443\u0435\u0442 \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u044c \u043d\u0430\u0445\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u0432 \u0442\u0435\u043b\u0435 \u043c\u0435\u0442\u043e\u0434\u0430 <em>createAndRegisterListener(String topic).\u00a0<\/em>\u041b\u043e\u0433\u0438\u043a\u0430 \u043d\u0430\u0445\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u0432 \u0442\u043e\u043c \u0436\u0435 \u043a\u043b\u0430\u0441\u0441\u0435<em>\u00a0KafkaListenerCreator.<\/em><\/p>\n<pre><code class=\"java\">@Service public class KafkaListenerCreator {   \/\/... HERE HAS TO BE VARIABLES FROM PREVIOUS EXAMPLE    @Autowired   private KafkaListenerEndpointRegistry kafkaListenerEndpointRegistry;   @Autowired   private KafkaListenerContainerFactory kafkaListenerContainerFactory;    public void createAndRegisterListener(String topic) {     KafkaListenerEndpoint listener = createKafkaListenerEndpoint(topic);     kafkaListenerEndpointRegistry.registerListenerContainer(listener, kafkaListenerContainerFactory, true);   }    \/\/... HERE HAS TO BE METHODS FROM PREVIOUS EXAMPLE  }<\/code><\/pre>\n<h3>4. \u0421\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u043e\u043a\u0440\u0443\u0436\u0435\u043d\u0438\u0435 \u0434\u043b\u044f \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f<\/h3>\n<p>\u0421\u043d\u0430\u0447\u0430\u043b\u0430 \u0441\u043e\u0437\u0434\u0430\u0434\u0438\u043c REST \u043c\u0435\u0442\u043e\u0434 \u0434\u043b\u044f \u0434\u043b\u044f \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u044f \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u044f Kafka.<\/p>\n<p>\u042f \u0441\u043e\u0437\u0434\u0430\u043b \u043a\u043b\u0430\u0441\u0441 <em>KafkaController \u0438 \u043c\u0435\u0442\u043e\u0434 create(String topic). <\/em>\u042d\u0442\u043e\u0442 \u043c\u0435\u0442\u043e\u0434 \u043c\u043e\u0436\u0435\u0442 \u0431\u044b\u0442\u044c \u0432\u044b\u0437\u0432\u0430\u043d \u0441 \u043f\u043e\u043c\u043e\u0449\u044c\u044e POST HTTP<em> \u0437\u0430\u043f\u0440\u043e\u0441\u0430.<\/em><\/p>\n<pre><code class=\"java\">@RestController public class KafkaController {     @Autowired     KafkaListenerCreator kafkaListenerCreator;      @PostMapping(path = \"\/create\")     @ResponseStatus(HttpStatus.OK)     public void create(@RequestParam String topic) {         kafkaListenerCreator.createAndRegisterListener(topic);     } }<\/code><\/pre>\n<p>\u0414\u0430\u043b\u0435\u0435 \u0441\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u043c\u0435\u0442\u043e\u0434 \u0434\u043b\u044f \u043e\u0442\u043f\u0440\u0430\u0432\u043a\u0438 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f \u0432 \u043d\u043e\u0432\u043e \u0441\u043e\u0437\u0434\u0430\u043d\u043d\u044b\u0439 \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u044c.<\/p>\n<p>\u041c\u0435\u0442\u043e\u0434\u00a0<em>send(@RequestParam String topic, @RequestParam String message)\u00a0<\/em>\u0438\u043c\u0435\u0435\u0442 \u0434\u0432\u0430 \u043f\u0430\u0440\u0430\u043c\u0435\u0442\u0440\u0430. \u0413\u0434\u0435\u00a0<em>\u201ctopic\u201d<\/em>\u00a0\u044d\u0442\u043e \u0438\u043c\u044f \u0442\u043e\u043f\u0438\u043a\u0430 \u0443 \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u044f Kafka, \u0430\u00a0<em>\u201cmessage\u201d<\/em>\u00a0\u044d\u0442\u043e \u0442\u0435\u043a\u0441\u0442 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f, \u043a\u043e\u0442\u043e\u0440\u043e\u0435 \u043e\u0442\u043f\u0440\u0430\u0432\u043b\u044f\u0435\u0442\u0441\u044f \u0447\u0435\u0440\u0435\u0437 <em>KafkaTemplate<\/em>.<\/p>\n<pre><code class=\"java\">@RestController public class KafkaController {     \/\/... HERE HAS TO BE VARIABLES FROM PREVIOUS EXAMPLE     @Autowired     private KafkaTemplate&lt;String, String> kafkaTemplate;      @PostMapping(path = \"\/send\")     @ResponseStatus(HttpStatus.OK)     public void send(@RequestParam String topic, @RequestParam String message) {         kafkaTemplate.send(topic, message);     }     \/\/... HERE HAS TO BE METHODS FROM PREVIOUS EXAMPLE }<\/code><\/pre>\n<h3>5. \u041f\u0440\u043e\u0442\u0435\u0441\u0442\u0438\u0440\u0443\u0435\u043c \u0440\u0435\u0448\u0435\u043d\u0438\u0435<\/h3>\n<ol>\n<li>\n<p>\u0432\u044b\u0437\u0432\u0430\u0442\u044c http:\/\/localhost:8080\/create?topic=myTopic1<\/p>\n<\/li>\n<li>\n<p>\u0432\u044b\u0437\u0432\u0430\u0442\u044c http:\/\/localhost:8080\/send?topic=myTopic1&amp;message=myTxt1<\/p>\n<\/li>\n<li>\n<p>\u043e\u0436\u0438\u0434\u0430\u0435\u043c\u044b\u0439 \u043b\u043e\u0433: \u201c RECORD PROCESSING: myTxt1\u201d<\/p>\n<\/li>\n<\/ol>\n<h3>6. \u0417\u0430\u043a\u043b\u044e\u0447\u0435\u043d\u0438\u0435<\/h3>\n<p>\u042d\u0442\u0430 \u0441\u0442\u0430\u0442\u044c\u044f \u043f\u0440\u0435\u0434\u043e\u0441\u0442\u0430\u0432\u043b\u044f\u0435\u0442 \u0431\u044b\u0441\u0442\u0440\u043e\u0435 \u0440\u0435\u0448\u0435\u043d\u0438\u0435 \u0434\u043b\u044f \u043f\u0440\u043e\u0431\u043b\u0435\u043c\u044b \u0434\u0438\u043d\u0430\u043c\u0438\u0447\u0435\u0441\u043a\u043e\u0433\u043e \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u044f \u0441\u043b\u0443\u0448\u0430\u0442\u0435\u043b\u0435\u0439 \u0432 Kafka.<\/p>\n<h2>\u0414\u043e\u043f\u043e\u043b\u043d\u0438\u0442\u0435\u043b\u044c\u043d\u044b\u0435 \u0440\u0435\u0441\u0443\u0440\u0441\u044b<\/h2>\n<p><a href=\"https:\/\/medium.com\/bliblidotcom-techblog\/dynamic-spring-boot-kafka-consumer-af8740f2c703?source=post_page-----4f8f359d715e--------------------------------\" rel=\"noopener noreferrer nofollow\">https:\/\/medium.com\/bliblidotcom-techblog\/dynamic-spring-boot-kafka-consumer-af8740f2c703?source=post_page&#8212;&#8212;4f8f359d715e&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;<\/a><\/p>\n<p><a href=\"https:\/\/github.com\/spring-projects\/spring-kafka\/issues\/460?source=post_page-----4f8f359d715e--------------------------------\" rel=\"noopener noreferrer nofollow\">https:\/\/github.com\/spring-projects\/spring-kafka\/issues\/460?source=post_page&#8212;&#8212;4f8f359d715e&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;&#8212;<\/a><\/p>\n<\/p>\n<\/div>\n<\/div>\n<\/div>\n<p><!----><!----><\/div>\n<p><!----><!----><br \/> \u0441\u0441\u044b\u043b\u043a\u0430 \u043d\u0430 \u043e\u0440\u0438\u0433\u0438\u043d\u0430\u043b \u0441\u0442\u0430\u0442\u044c\u0438 <a href=\"https:\/\/habr.com\/ru\/articles\/758958\/\"> https:\/\/habr.com\/ru\/articles\/758958\/<\/a><br \/><\/br><\/br><\/p>\n","protected":false},"author":1,"featured_media":0,"comment_status":"open","ping_status":"open","sticky":false,"template":"","format":"standard","meta":{"footnotes":""},"categories":[],"tags":[],"class_list":["post-354025","post","type-post","status-publish","format-standard","hentry"],"_links":{"self":[{"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=\/wp\/v2\/posts\/354025","targetHints":{"allow":["GET"]}}],"collection":[{"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=\/wp\/v2\/posts"}],"about":[{"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=\/wp\/v2\/types\/post"}],"author":[{"embeddable":true,"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=\/wp\/v2\/users\/1"}],"replies":[{"embeddable":true,"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=%2Fwp%2Fv2%2Fcomments&post=354025"}],"version-history":[{"count":0,"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=\/wp\/v2\/posts\/354025\/revisions"}],"wp:attachment":[{"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=%2Fwp%2Fv2%2Fmedia&parent=354025"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=%2Fwp%2Fv2%2Fcategories&post=354025"},{"taxonomy":"post_tag","embeddable":true,"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=%2Fwp%2Fv2%2Ftags&post=354025"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}