{"id":378097,"date":"2024-06-05T15:00:28","date_gmt":"2024-06-05T15:00:28","guid":{"rendered":"http:\/\/savepearlharbor.com\/?p=378097"},"modified":"-0001-11-30T00:00:00","modified_gmt":"-0001-11-29T21:00:00","slug":"","status":"publish","type":"post","link":"https:\/\/savepearlharbor.com\/?p=378097","title":{"rendered":"<span>E2E-\u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435 Flink Job \u0441 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>\u041f\u0440\u0438\u0432\u0435\u0442, \u0425\u0430\u0431\u0440! \u0421 \u0432\u0430\u043c\u0438 \u0410\u043b\u0435\u043a\u0441\u0430\u043d\u0434\u0440 \u0411\u043e\u0431\u0440\u044f\u043a\u043e\u0432, \u0442\u0435\u0445\u043b\u0438\u0434 \u0432 \u043a\u043e\u043c\u0430\u043d\u0434\u0435 \u041c\u0422\u0421 \u0410\u043d\u0430\u043b\u0438\u0442\u0438\u043a\u0438. \u042f \u043a \u0432\u0430\u043c \u0441 \u043d\u043e\u0432\u043e\u0439 \u0441\u0442\u0430\u0442\u044c\u0451\u0439 \u0438\u0437 \u0446\u0438\u043a\u043b\u0430 \u043f\u0440\u043e \u0444\u0440\u0435\u0439\u043c\u0432\u043e\u0440\u043a Apache Flink. <\/p>\n<p>\u0412 <a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/812905\/\">\u043f\u0440\u0435\u0434\u044b\u0434\u0443\u0449\u0435\u0439 \u0447\u0430\u0441\u0442\u0438<\/a> \u044f \u0440\u0430\u0441\u0441\u043a\u0430\u0437\u0430\u043b, \u043a\u0430\u043a \u0441\u043e\u0437\u0434\u0430\u0442\u044c Unit-\u0442\u0435\u0441\u0442 \u043d\u0430 \u043f\u043e\u043b\u043d\u043e\u0446\u0435\u043d\u043d\u0443\u044e \u0434\u0436\u043e\u0431\u0443 Flink \u0438 \u043e\u0442\u0434\u0435\u043b\u044c\u043d\u044b\u0435 stateful-\u043e\u043f\u0435\u0440\u0430\u0442\u043e\u0440\u044b \u0441 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u043d\u0438\u0435\u043c Flink MiniCluster. \u0415\u0449\u0451 \u043c\u044b \u043d\u0430\u0443\u0447\u0438\u043b\u0438\u0441\u044c \u0437\u0430\u043f\u0443\u0441\u043a\u0430\u0442\u044c \u043c\u0438\u043d\u0438-\u043a\u043b\u0430\u0441\u0442\u0435\u0440 \u043e\u0434\u0438\u043d \u0440\u0430\u0437 \u043f\u0435\u0440\u0435\u0434 \u0432\u0441\u0435\u043c\u0438 \u0442\u0435\u0441\u0442\u043e\u0432\u044b\u043c\u0438 \u043a\u043b\u0430\u0441\u0441\u0430\u043c\u0438, \u043a\u043e\u0442\u043e\u0440\u044b\u0435 \u043d\u0443\u0436\u0434\u0430\u044e\u0442\u0441\u044f \u0432 \u043d\u0451\u043c. \u0412 \u0434\u043e\u043f\u043e\u043b\u043d\u0435\u043d\u0438\u0435 \u0441\u043e\u0437\u0434\u0430\u043b\u0438 \u0432\u0441\u043f\u043e\u043c\u043e\u0433\u0430\u0442\u0435\u043b\u044c\u043d\u044b\u0435 \u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u0438 \u0438 \u0430\u043d\u043d\u043e\u0442\u0430\u0446\u0438\u0438, \u0437\u043d\u0430\u0447\u0438\u0442\u0435\u043b\u044c\u043d\u043e \u0440\u0430\u0437\u0434\u0435\u043b\u044f\u044f \u043e\u0442\u0432\u0435\u0442\u0441\u0442\u0432\u0435\u043d\u043d\u043e\u0441\u0442\u044c \u0432 \u0442\u0435\u0441\u0442\u0430\u0445 \u0438 \u0443\u043f\u0440\u043e\u0449\u0430\u044f \u043b\u043e\u0433\u0438\u043a\u0443 \u043d\u0430\u043f\u0438\u0441\u0430\u043d\u0438\u044f \u043d\u043e\u0432\u044b\u0445 \u0442\u0435\u0441\u0442\u043e\u0432. <\/p>\n<p>\u0412 \u043f\u0440\u0435\u0434\u044b\u0434\u0443\u0449\u0438\u0445 \u0442\u0435\u0441\u0442\u0430\u0445 \u043d\u0430 \u0434\u0436\u043e\u0431\u0443 \u043c\u044b \u043d\u0435 \u0437\u0430\u0442\u0440\u0430\u0433\u0438\u0432\u0430\u043b\u0438 \u0438\u043d\u0442\u0435\u0433\u0440\u0430\u0446\u0438\u044e \u0441 Kafka, \u0432\u0435\u0434\u044c \u043d\u0430\u043c \u0431\u044b\u043b\u0438 \u043d\u0435 \u0432\u0430\u0436\u043d\u044b \u0440\u0435\u0430\u043b\u044c\u043d\u044b\u0435 source \u0438 sink. \u0412 \u044d\u0442\u043e\u0439 \u0441\u0442\u0430\u0442\u044c\u0435 \u043f\u0440\u043e\u0434\u043e\u043b\u0436\u0438\u043c \u0440\u0430\u0437\u0431\u0438\u0440\u0430\u0442\u044c\u0441\u044f \u0432 \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0438 \u0438 \u043d\u0430\u043f\u0438\u0448\u0435\u043c \u043f\u043e\u043b\u043d\u043e\u0446\u0435\u043d\u043d\u044b\u0439 E2E-\u0442\u0435\u0441\u0442, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043e\u0445\u0432\u0430\u0442\u0438\u0442 Kafka \u0438 Flink \u0432\u043c\u0435\u0441\u0442\u0435 \u0441 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u043d\u0438\u0435\u043c Testcontainers. \u0422\u0430\u043a\u0436\u0435 \u0440\u0430\u0441\u0441\u043c\u043e\u0442\u0440\u0438\u043c \u043d\u0435\u043e\u0447\u0435\u0432\u0438\u0434\u043d\u044b\u0435 \u043f\u0440\u043e\u0431\u043b\u0435\u043c\u044b \u0432 \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0438 \u0438 \u043d\u043e\u0432\u044b\u0435 \u0443\u043d\u0438\u0432\u0435\u0440\u0441\u0430\u043b\u044c\u043d\u044b\u0435 \u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u0438.<\/p>\n<figure class=\"full-width\"><img loading=\"lazy\" decoding=\"async\" src=\"https:\/\/habrastorage.org\/r\/w780q1\/getpro\/habr\/upload_files\/e61\/644\/86e\/e6164486e65ca746f90ea6d234daee7a.jpg\" width=\"1280\" height=\"720\" data-src=\"https:\/\/habrastorage.org\/getpro\/habr\/upload_files\/e61\/644\/86e\/e6164486e65ca746f90ea6d234daee7a.jpg\" data-blurred=\"true\"\/><\/figure>\n<details class=\"spoiler\">\n<summary>\u0421\u043f\u0438\u0441\u043e\u043a \u043c\u043e\u0438\u0445 \u043f\u043e\u0441\u0442\u043e\u0432 \u043f\u0440\u043e Flink<\/summary>\n<div class=\"spoiler__content\">\n<ol>\n<li>\n<p><a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/772898\/\">\u0412\u0432\u0435\u0434\u0435\u043d\u0438\u0435 \u0432 Apache Flink: \u043e\u0441\u0432\u0430\u0438\u0432\u0430\u0435\u043c \u0444\u0440\u0435\u0439\u043c\u0432\u043e\u0440\u043a \u043d\u0430 \u0440\u0435\u0430\u043b\u044c\u043d\u044b\u0445 \u043f\u0440\u0438\u043c\u0435\u0440\u0430\u0445<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/775970\/\">\u041f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0435 \u043f\u043e\u0434 Apache Flink: \u0441 \u0447\u0435\u0433\u043e \u043d\u0430\u0447\u0430\u0442\u044c?<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/786012\/\">Apache Flink. \u041a\u0430\u043a \u0440\u0430\u0431\u043e\u0442\u0430\u0435\u0442 \u0434\u0435\u0434\u0443\u043f\u043b\u0438\u043a\u0430\u0446\u0438\u044f \u0434\u0430\u043d\u043d\u044b\u0445 \u0432 \u043f\u043e\u0442\u043e\u043a\u0435 Kafka-to-Kafka?<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/798667\/\">\u0414\u0438\u043d\u0430\u043c\u0438\u0447\u0435\u0441\u043a\u043e\u0435 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u0438\u0435 \u0432\u044b\u0445\u043e\u0434\u043d\u043e\u0433\u043e \u0442\u043e\u043f\u0438\u043a\u0430<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/801693\/\">\u041a\u0430\u043a \u043f\u0440\u043e\u0432\u0435\u0441\u0442\u0438 unit-\u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435 Flink-\u043e\u043f\u0435\u0440\u0430\u0442\u043e\u0440\u043e\u0432: Test Harness<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/812905\/\">Unit-\u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435 Flink-\u043e\u043f\u0435\u0440\u0430\u0442\u043e\u0440\u043e\u0432, Job: Flink MiniCluster<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/819681\/\"><strong>E2E-\u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435 Flink Job \u0441 Kafka<\/strong><\/a><\/p>\n<\/li>\n<\/ol>\n<\/div>\n<\/details>\n<p>\u0412\u0435\u0441\u044c \u0440\u0430\u0437\u0431\u0438\u0440\u0430\u0435\u043c\u044b\u0439 \u0438\u0441\u0445\u043e\u0434\u043d\u044b\u0439 \u043a\u043e\u0434 \u043c\u043e\u0436\u043d\u043e \u043d\u0430\u0439\u0442\u0438 \u0432 \u0440\u0435\u043f\u043e\u0437\u0438\u0442\u043e\u0440\u0438\u0438 <a href=\"https:\/\/github.com\/AlexanderBobryakov\/flink-spring\">AlexanderBobryakov\/flink-spring<\/a>. \u0412 master-\u0432\u0435\u0442\u043a\u0435 \u2014 \u0438\u0442\u043e\u0433\u043e\u0432\u044b\u0439 \u043f\u0440\u043e\u0435\u043a\u0442 \u043f\u043e \u0432\u0441\u0435\u0439 \u0441\u0435\u0440\u0438\u0438 \u0441\u0442\u0430\u0442\u0435\u0439. \u042d\u0442\u0430 \u0441\u0442\u0430\u0442\u044c\u044f \u0441\u043e\u043e\u0442\u0432\u0435\u0442\u0441\u0442\u0432\u0443\u0435\u0442 \u0440\u0435\u043b\u0438\u0437\u043d\u043e\u0439 \u0432\u0435\u0442\u043a\u0435 \u043f\u043e\u0434 \u043d\u0430\u0437\u0432\u0430\u043d\u0438\u0435\u043c <a href=\"https:\/\/github.com\/AlexanderBobryakov\/flink-spring\/tree\/release\/6_e2e_deduplicator_test\">release\/6_e2e_deduplicator_test<\/a>.<\/p>\n<details class=\"spoiler\">\n<summary>\u041e\u0433\u043b\u0430\u0432\u043b\u0435\u043d\u0438\u0435 \u0441\u0442\u0430\u0442\u044c\u0438<\/summary>\n<div class=\"spoiler__content\">\n<ol>\n<li>\n<p><a href=\"#1\">E2E-\u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#2\">\u041f\u043e\u0434\u043d\u0438\u043c\u0430\u0435\u043c Kafka \u0441 \u043f\u043e\u043c\u043e\u0449\u044c\u044e Testcontainers<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#3\">\u0410\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u0438 \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u0434\u043b\u044f Kafka<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#4\">KafkaTestConsumer<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#5\">TestKafkaFacade<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#6\">KafkaTopicCreatorConfig<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#7\">E2E-\u0442\u0435\u0441\u0442 \u043d\u0430 Flink Job<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#8\">\u0411\u0435\u0437\u043e\u043f\u0430\u0441\u043d\u043e\u0435 \u0437\u0430\u0432\u0435\u0440\u0448\u0435\u043d\u0438\u0435 E2E-\u0442\u0435\u0441\u0442\u043e\u0432<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#9\">RocksDB \u0432 E2E-\u0442\u0435\u0441\u0442\u0430\u0445<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#10\">\u0412\u044b\u0432\u043e\u0434<\/a><\/p>\n<\/li>\n<\/ol>\n<\/div>\n<\/details>\n<p><a class=\"anchor\" name=\"1\" id=\"1\"><\/a><\/p>\n<h2>E2E-\u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435<\/h2>\n<p>E2E-\u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435 \u043e\u0445\u0432\u0430\u0442\u044b\u0432\u0430\u0435\u0442 \u043f\u043e\u0432\u0435\u0434\u0435\u043d\u0438\u0435 \u0432\u0441\u0435\u0439 \u0441\u0438\u0441\u0442\u0435\u043c\u044b \u043e\u0442 \u043d\u0430\u0447\u0430\u043b\u0430 \u0434\u043e \u043a\u043e\u043d\u0446\u0430. \u041f\u043e \u0441\u043f\u0435\u0446\u0438\u0444\u0438\u043a\u0435 \u043d\u0430\u0448\u0435\u0439 \u0434\u0436\u043e\u0431\u044b, \u043a\u043e\u0442\u043e\u0440\u0443\u044e \u043c\u044b \u0440\u0430\u0441\u0441\u043c\u0430\u0442\u0440\u0438\u0432\u0430\u043b\u0438 \u0432 \u043f\u0440\u0435\u0434\u044b\u0434\u0443\u0449\u0438\u0445 \u0447\u0430\u0441\u0442\u044f\u0445, \u0432 \u043d\u0430\u0447\u0430\u043b\u0435 \u0435\u0441\u0442\u044c Kafka-\u0442\u043e\u043f\u0438\u043a \u0441 \u0434\u0430\u043d\u043d\u044b\u043c\u0438 ClickMessage, \u0430 \u043d\u0430 \u0432\u044b\u0445\u043e\u0434\u0435 \u2014 \u043c\u043d\u043e\u0433\u043e \u0440\u0430\u0437\u043d\u044b\u0445 product-\u0442\u043e\u043f\u0438\u043a\u043e\u0432. \u0417\u043d\u0430\u0447\u0438\u0442, \u0432 \u0442\u0435\u0441\u0442\u0435 \u0432\u0441\u0451 \u044d\u0442\u043e \u0434\u043e\u043b\u0436\u043d\u043e \u0443\u0447\u0438\u0442\u044b\u0432\u0430\u0442\u044c\u0441\u044f. <\/p>\n<p>\u0426\u0435\u043b\u044c \u0442\u0430\u043a\u043e\u0433\u043e \u0442\u0435\u0441\u0442\u0430 \u2014 \u043f\u0440\u043e\u0432\u0435\u0440\u0438\u0442\u044c \u0432\u0441\u0435 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u043c\u044b\u0435 \u0438\u043d\u0442\u0435\u0433\u0440\u0430\u0446\u0438\u0438 \u043c\u0435\u0436\u0434\u0443 \u0441\u043e\u0431\u043e\u0439. \u0418\u043d\u0430\u0447\u0435 \u0431\u043b\u043e\u043a\u0438 \u0441\u0438\u0441\u0442\u0435\u043c\u044b \u043c\u043e\u0433\u0443\u0442 \u0440\u0430\u0431\u043e\u0442\u0430\u0442\u044c \u043f\u043e \u043e\u0442\u0434\u0435\u043b\u044c\u043d\u043e\u0441\u0442\u0438, \u0430 \u0432\u043c\u0435\u0441\u0442\u0435 \u0432\u0441\u0451 \u0441\u043b\u043e\u043c\u0430\u0435\u0442\u0441\u044f. \u0412 \u043d\u0430\u0448\u0435\u043c \u0441\u043b\u0443\u0447\u0430\u0435 \u0434\u043e\u043b\u0436\u0435\u043d \u043f\u043e\u0434\u043d\u0438\u043c\u0430\u0442\u044c\u0441\u044f \u0432\u0435\u0441\u044c Spring-\u043a\u043e\u043d\u0442\u0435\u043a\u0441\u0442, \u0441\u0442\u0430\u0440\u0442\u043e\u0432\u0430\u0442\u044c Flink \u0438 Kafka, \u043a\u0430\u043a \u0431\u0443\u0434\u0442\u043e \u043c\u044b \u0437\u0430\u043f\u0443\u0441\u043a\u0430\u0435\u043c \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0435 \u043d\u0430 \u043f\u0440\u043e\u0434\u0435.<\/p>\n<p><a class=\"anchor\" name=\"2\" id=\"2\"><\/a><\/p>\n<h2>\u041f\u043e\u0434\u043d\u0438\u043c\u0430\u0435\u043c Kafka \u0441 \u043f\u043e\u043c\u043e\u0449\u044c\u044e Testcontainers<\/h2>\n<p><a href=\"https:\/\/testcontainers.com\">Testcontainers<\/a> \u2014 \u200b\u200b\u044d\u0442\u043e \u0431\u0438\u0431\u043b\u0438\u043e\u0442\u0435\u043a\u0430 Java, \u043a\u043e\u0442\u043e\u0440\u0430\u044f \u043f\u043e\u0434\u0434\u0435\u0440\u0436\u0438\u0432\u0430\u0435\u0442 \u0442\u0435\u0441\u0442\u044b JUnit. \u041e\u043d\u0430 \u0434\u0430\u0451\u0442 \u0432\u043e\u0437\u043c\u043e\u0436\u043d\u043e\u0441\u0442\u044c \u0437\u0430\u043f\u0443\u0441\u043a\u0430\u0442\u044c \u0432 \u043d\u0438\u0445 \u0432\u0441\u0451, \u0447\u0442\u043e \u043c\u043e\u0436\u0435\u0442 \u0437\u0430\u043f\u0443\u0441\u043a\u0430\u0442\u044c\u0441\u044f \u0432 Docker. \u0417\u043d\u0430\u0447\u0438\u0442, \u0432\u044b \u043c\u043e\u0436\u0435\u0442\u0435 \u043f\u0440\u043e\u0432\u0435\u0440\u0438\u0442\u044c \u043b\u044e\u0431\u0443\u044e \u0438\u043d\u0442\u0435\u0433\u0440\u0430\u0446\u0438\u044e \u0432\u0430\u0448\u0435\u0433\u043e \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u044f: \u0441 \u0411\u0414, \u0431\u0440\u043e\u043a\u0435\u0440\u0430\u043c\u0438 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0439, \u0434\u0440\u0443\u0433\u0438\u043c\u0438 \u0441\u0435\u0440\u0432\u0438\u0441\u0430\u043c\u0438 \u0438 \u0442\u0430\u043a \u0434\u0430\u043b\u0435\u0435. <\/p>\n<p>\u0421\u0446\u0435\u043d\u0430\u0440\u0438\u0439 \u043d\u0430\u043f\u0438\u0441\u0430\u043d\u0438\u044f \u0442\u0435\u0441\u0442\u0430 \u0432 \u0438\u0442\u043e\u0433\u0435 \u0432\u044b\u0433\u043b\u044f\u0434\u0438\u0442 \u0442\u0430\u043a:<\/p>\n<ol>\n<li>\n<p>\u041e\u043f\u0440\u0435\u0434\u0435\u043b\u0438\u0442\u044c \u0432 \u0442\u0435\u0441\u0442\u0435 Testcontainers \u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440 \u2014 \u043d\u0430\u043f\u0440\u0438\u043c\u0435\u0440, \u0434\u043b\u044f Kafka.<\/p>\n<\/li>\n<li>\n<p>\u0417\u0430\u043f\u0443\u0441\u0442\u0438\u0442\u044c Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440.<\/p>\n<\/li>\n<li>\n<p>\u041f\u0440\u043e\u0431\u0440\u043e\u0441\u0438\u0442\u044c \u0441\u0432\u043e\u0439\u0441\u0442\u0432\u0430 \u0434\u043b\u044f \u043f\u043e\u0434\u043a\u043b\u044e\u0447\u0435\u043d\u0438\u044f \u043a Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440\u0443 \u0432 \u043a\u043e\u043d\u0444\u0438\u0433 \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u044f.<\/p>\n<\/li>\n<li>\n<p>\u0417\u0430\u043f\u0443\u0441\u0442\u0438\u0442\u044c \u0442\u0435\u0441\u0442, \u0432 \u043a\u043e\u0442\u043e\u0440\u043e\u043c \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0435 \u043f\u043e\u0434\u043a\u043b\u044e\u0447\u0430\u0435\u0442\u0441\u044f \u0441\u043e\u0433\u043b\u0430\u0441\u043d\u043e \u043a\u043e\u043d\u0444\u0438\u0433\u0443 \u043a Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440\u0443.<\/p>\n<\/li>\n<\/ol>\n<p>\u0412 \u043f\u0440\u0435\u0434\u044b\u0434\u0443\u0449\u0435\u0439 \u0441\u0442\u0430\u0442\u044c\u0435 \u043c\u044b \u0437\u0430\u0442\u0440\u043e\u043d\u0443\u043b\u0438 \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u0435 \u043a\u0430\u0441\u0442\u043e\u043c\u043d\u044b\u0445 \u0430\u043d\u043d\u043e\u0442\u0430\u0446\u0438\u0439 \u043f\u043e\u0434 \u0442\u0435\u0441\u0442\u044b. \u042d\u0442\u043e \u0431\u044b\u043b\u043e \u0434\u043e\u0441\u0442\u0430\u0442\u043e\u0447\u043d\u043e \u0443\u0434\u043e\u0431\u043d\u043e, \u043f\u043e\u044d\u0442\u043e\u043c\u0443 \u043f\u0440\u0435\u0434\u043b\u0430\u0433\u0430\u044e \u043f\u0440\u0438\u0434\u0435\u0440\u0436\u0438\u0432\u0430\u0442\u044c\u0441\u044f \u0430\u043d\u0430\u043b\u043e\u0433\u0438\u0447\u043d\u043e\u0433\u043e \u043f\u043e\u0434\u0445\u043e\u0434\u0430 \u0438 \u0432 \u044d\u0442\u043e\u0442 \u0440\u0430\u0437. \u0414\u043b\u044f \u043d\u0430\u0447\u0430\u043b\u0430 \u043d\u0430\u043c \u043d\u0443\u0436\u043d\u043e \u043f\u043e\u0434\u043a\u043b\u044e\u0447\u0435\u043d\u0438\u0435 \u043a Kafka. TestContainers \u043f\u0440\u0435\u0434\u043e\u0441\u0442\u0430\u0432\u043b\u044f\u0435\u0442 Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440 <a href=\"https:\/\/www.testcontainers.org\/modules\/kafka\/\">\u00ab\u0438\u0437 \u043a\u043e\u0440\u043e\u0431\u043a\u0438\u00bb<\/a>, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043d\u0435 \u043d\u0443\u0436\u0434\u0430\u0435\u0442\u0441\u044f \u0432 <strong>\u043e\u0442\u0434\u0435\u043b\u044c\u043d\u043e\u0439<\/strong> \u0438\u043d\u0438\u0446\u0438\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438 Zookeeper. \u0415\u0449\u0451 TestContainers \u043c\u043e\u0436\u043d\u043e \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c, \u0447\u0442\u043e\u0431\u044b \u0441\u043e\u0437\u0434\u0430\u0442\u044c \u0432\u0441\u0435 \u043d\u0435\u043e\u0431\u0445\u043e\u0434\u0438\u043c\u044b\u0435 \u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440\u044b \u043d\u0430 \u043e\u0441\u043d\u043e\u0432\u0435 \u043b\u044e\u0431\u044b\u0445 \u0434\u043e\u043a\u0435\u0440-\u043e\u0431\u0440\u0430\u0437\u043e\u0432.<\/p>\n<p>\u0414\u043b\u044f \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u043d\u0438\u044f \u0433\u043e\u0442\u043e\u0432\u043e\u0433\u043e Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440\u0430 \u043c\u043e\u0436\u0435\u043c \u0432\u043e\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c\u0441\u044f \u0437\u0430\u0432\u0438\u0441\u0438\u043c\u043e\u0441\u0442\u044c\u044e:<\/p>\n<pre><code>testImplementation \"org.testcontainers:kafka\"<\/code><\/pre>\n<p>\u0410\u043d\u0430\u043b\u043e\u0433\u0438\u0447\u043d\u043e \u0442\u043e\u043c\u0443, \u043a\u0430\u043a \u043c\u044b \u0441\u043e\u0437\u0434\u0430\u043b\u0438 JUnit Extension \u0434\u043b\u044f \u0441\u0442\u0430\u0440\u0442\u0430 Flink MiniCluster, \u0441\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u043d\u043e\u0432\u044b\u0439 Extension \u0434\u043b\u044f \u0441\u0442\u0430\u0440\u0442\u0430 Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440\u0430:<\/p>\n<pre><code class=\"java\">@Slf4j @SuppressWarnings({\"PMD.AvoidUsingVolatile\"}) public class KafkaContainerExtension implements BeforeAllCallback, ExtensionContext.Store.CloseableResource {    private static final KafkaContainer KAFKA =        new KafkaContainer(DockerImageName.parse(\"confluentinc\/cp-kafka:7.3.2\"))            .withEnv(\"KAFKA_AUTO_CREATE_TOPICS_ENABLE\", \"false\");     private static final Lock LOCK = new ReentrantLock();    private static volatile boolean started;     @Override    public void beforeAll(ExtensionContext context) {        LOCK.lock();        try {            if (!started) {                log.info(\"Start Kafka Container\");                started = true;                Startables.deepStart(KAFKA).join();                System.setProperty(\"spring.kafka.bootstrap-servers\", KAFKA.getBootstrapServers());                System.setProperty(\"kafka.bootstrap-servers\", KAFKA.getBootstrapServers());                System.setProperty(\"spring.kafka.consumer.group-id\", \"group-id-spring\");                context.getRoot().getStore(GLOBAL).put(\"Kafka Container\", this);            }        } finally {            LOCK.unlock();        }    }     @Override    public void close() {        log.info(\"Close Kafka Container\");        KAFKA.close();        started = false;    } }<\/code><\/pre>\n<p>\u041a\u043e\u0434 \u043e\u0447\u0435\u043d\u044c \u043f\u043e\u0445\u043e\u0436 \u043d\u0430 \u043d\u0430\u0448 \u0441\u0443\u0449\u0435\u0441\u0442\u0432\u0443\u044e\u0449\u0438\u0439 FlinkClusterExtension, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u044f \u043e\u043f\u0438\u0441\u0430\u043b \u0432 \u043f\u0440\u043e\u0448\u043b\u043e\u0439 \u0441\u0442\u0430\u0442\u044c\u0435: \u043c\u044b \u0438\u043d\u0438\u0446\u0438\u0430\u043b\u0438\u0437\u0438\u0440\u0443\u0435\u043c Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440 \u043f\u043e \u0443\u043a\u0430\u0437\u0430\u043d\u043d\u043e\u043c\u0443 \u0434\u043e\u043a\u0435\u0440-\u043e\u0431\u0440\u0430\u0437\u0443, \u043f\u043e\u0442\u043e\u043c \u0432 beforeAll() \u0441\u0438\u043d\u0445\u0440\u043e\u043d\u043d\u043e \u0437\u0430\u043f\u0443\u0441\u043a\u0430\u0435\u043c \u0435\u0433\u043e \u0447\u0435\u0440\u0435\u0437 \u0432\u044b\u0437\u043e\u0432 Startables.deepStart(KAFKA).join(), \u043e\u0431\u0432\u044f\u0437\u044b\u0432\u0430\u044f \u0431\u043b\u043e\u043a\u0438\u0440\u043e\u0432\u043a\u0430\u043c\u0438. \u0412 \u043a\u043e\u043d\u0446\u0435 \u0432\u044b\u043f\u043e\u043b\u043d\u0435\u043d\u0438\u044f \u0432\u0441\u0435\u0445 \u0437\u0430\u0432\u0438\u0441\u0438\u043c\u044b\u0445 \u0442\u0435\u0441\u0442\u043e\u0432 \u0437\u0430\u043a\u0440\u044b\u0432\u0430\u0435\u043c \u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440 \u0432 \u043a\u043e\u043b\u0431\u044d\u043a-\u043c\u0435\u0442\u043e\u0434\u0435 \u0436\u0438\u0437\u043d\u0435\u043d\u043d\u043e\u0433\u043e \u0446\u0438\u043a\u043b\u0430 JUnit \u0442\u0435\u0441\u0442\u043e\u0432 close().<\/p>\n<p>\u0412\u043e\u0437\u043d\u0438\u043a\u0430\u0435\u0442 \u0432\u043e\u043f\u0440\u043e\u0441: \u043a\u0430\u043a \u043d\u0430\u0448\u0435 Spring-\u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0435 \u043f\u0440\u0438 \u0441\u0442\u0430\u0440\u0442\u0435 \u0431\u0443\u0434\u0435\u0442 \u043f\u043e\u0434\u043a\u043b\u044e\u0447\u0430\u0442\u044c\u0441\u044f \u043a Kafka? \u0412\u0435\u0434\u044c \u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440 \u0437\u0430\u043f\u0443\u0441\u043a\u0430\u0435\u0442\u0441\u044f \u043d\u0430 \u0441\u043b\u0443\u0447\u0430\u0439\u043d\u043e\u043c \u0441\u0432\u043e\u0431\u043e\u0434\u043d\u043e\u043c \u043f\u043e\u0440\u0442\u0443. \u0414\u043b\u044f \u044d\u0442\u043e\u0433\u043e \u043c\u044b \u043f\u0435\u0440\u0435\u0434\u0430\u0451\u043c \u043d\u0430\u0441\u0442\u0440\u043e\u0439\u043a\u0438 \u0432 \u043f\u0435\u0440\u0435\u043c\u0435\u043d\u043d\u044b\u0435 \u043e\u043a\u0440\u0443\u0436\u0435\u043d\u0438\u044f \u0447\u0435\u0440\u0435\u0437 System.setProperty() \u0432 \u0441\u0442\u0430\u0442\u0438\u0447\u0435\u0441\u043a\u043e\u043c \u043a\u043e\u043d\u0442\u0435\u043a\u0441\u0442\u0435 \u043d\u0435\u043f\u043e\u0441\u0440\u0435\u0434\u0441\u0442\u0432\u0435\u043d\u043d\u043e \u043f\u0435\u0440\u0435\u0434 \u0441\u0442\u0430\u0440\u0442\u043e\u043c \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u044f. \u041d\u0430\u0441\u0442\u0440\u043e\u0439\u043a\u0438 \u043f\u0435\u0440\u0435\u0434\u0430\u0451\u043c \u0441\u043e\u0433\u043b\u0430\u0441\u043d\u043e \u0441\u0442\u0440\u0443\u043a\u0442\u0443\u0440\u0435 application.yml, \u043f\u043e\u0442\u043e\u043c\u0443 \u0447\u0442\u043e \u043e\u043d\u0438 \u0431\u0443\u0434\u0443\u0442 \u00ab\u043f\u0435\u0440\u0435\u0437\u0430\u0442\u0438\u0440\u0430\u0442\u044c\u0441\u044f\u00bb \u0438\u0437 \u0443\u043a\u0430\u0437\u0430\u043d\u043d\u044b\u0445 \u043f\u0435\u0440\u0435\u043c\u0435\u043d\u043d\u044b\u0445 \u043e\u043a\u0440\u0443\u0436\u0435\u043d\u0438\u044f \u2014 <a href=\"https:\/\/docs.spring.io\/spring-boot\/docs\/current\/reference\/html\/features.html#features.external-config.typesafe-configuration-properties.relaxed-binding.environment-variables\">\u0434\u043e\u043a\u0443\u043c\u0435\u043d\u0442\u0430\u0446\u0438\u044f<\/a>:<\/p>\n<pre><code class=\"java\">kafka:  group-id: group_id  bootstrap-servers: localhost:29092<\/code><\/pre>\n<p>\u0414\u043b\u044f Spring \u043f\u0435\u0440\u0435\u0434\u0430\u0451\u043c \u0441\u0432\u043e\u0439\u0441\u0442\u0432\u0430, \u0447\u0442\u043e\u0431\u044b \u0430\u0432\u0442\u043e\u043c\u0430\u0442\u0438\u0447\u0435\u0441\u043a\u0438 \u0441\u043e\u0437\u0434\u0430\u043b\u0438\u0441\u044c \u043e\u0441\u043d\u043e\u0432\u043d\u044b\u0435 Spring-\u0431\u0438\u043d\u044b \u0434\u043b\u044f \u0438\u043d\u0442\u0435\u0433\u0440\u0430\u0446\u0438\u0438 \u0441 Kafka. \u041d\u0430\u043f\u0440\u0438\u043c\u0435\u0440, KafkaTemplate \u2014 \u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u044f \u043d\u0430\u0434 Kafka Producer, \u043a\u043e\u0442\u043e\u0440\u0430\u044f \u0443\u043c\u0435\u0435\u0442 \u043e\u0442\u043f\u0440\u0430\u0432\u043b\u044f\u0442\u044c \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f \u0432 \u0442\u043e\u043f\u0438\u043a. \u041e\u043d \u043f\u043e\u043d\u0430\u0434\u043e\u0431\u0438\u0442\u0441\u044f \u043d\u0430\u043c \u0432 \u0440\u0430\u043c\u043a\u0430\u0445 \u0442\u0435\u0441\u0442\u043e\u0432, \u043f\u043e\u044d\u0442\u043e\u043c\u0443 \u0434\u043e\u0431\u0430\u0432\u0438\u043c \u0437\u0430\u0432\u0438\u0441\u0438\u043c\u043e\u0441\u0442\u044c \u0432 \u0442\u0435\u0441\u0442\u0430\u0445:<\/p>\n<pre><code class=\"java\">testImplementation \"org.springframework.kafka:spring-kafka\"<\/code><\/pre>\n<p>\u042d\u0442\u043e\u0433\u043e \u043c\u043e\u0436\u043d\u043e \u0431\u044b\u043b\u043e \u0431\u044b \u0434\u043e\u0431\u0438\u0442\u044c\u0441\u044f \u0438 \u0430\u043b\u044c\u0442\u0435\u0440\u043d\u0430\u0442\u0438\u0432\u043d\u044b\u043c \u0441\u043f\u043e\u0441\u043e\u0431\u043e\u043c: \u0447\u0435\u0440\u0435\u0437 Spring-\u0438\u043d\u0438\u0446\u0438\u0430\u043b\u0438\u0437\u0430\u0442\u043e\u0440\u044b. \u041d\u0430\u043f\u0440\u0438\u043c\u0435\u0440, \u043c\u043e\u0436\u043d\u043e \u0431\u044b\u043b\u043e \u0431\u044b \u0441\u043e\u0437\u0434\u0430\u0442\u044c \u0441\u0432\u043e\u044e \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u044e \u0438\u043d\u0442\u0435\u0440\u0444\u0435\u0439\u0441\u0430 ApplicationContextInitializer. \u041d\u043e \u043c\u044b \u043d\u0435 \u0431\u0443\u0434\u0435\u043c \u0440\u0430\u0437\u0431\u0438\u0440\u0430\u0442\u044c \u044d\u0442\u043e\u0442 \u0441\u043f\u043e\u0441\u043e\u0431 \u0432 \u0441\u0442\u0430\u0442\u044c\u0435, \u043f\u043e\u0442\u043e\u043c\u0443 \u0447\u0442\u043e \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u044f \u0447\u0435\u0440\u0435\u0437 Extension \u0432\u044b\u0433\u043b\u044f\u0434\u0438\u0442 \u043a\u0440\u0430\u0441\u0438\u0432\u0435\u0435.<\/p>\n<p>\u0412 \u0438\u0442\u043e\u0433\u0435, \u0447\u0442\u043e\u0431\u044b \u043f\u0440\u0438\u043c\u0435\u043d\u044f\u0442\u044c \u0442\u0435\u043a\u0443\u0449\u0438\u0439 Extension \u0432 \u0442\u0435\u0441\u0442\u0430\u0445, \u043d\u0430\u043c \u0431\u044b\u043b\u0430 \u0431\u044b \u043f\u043e\u043b\u0435\u0437\u043d\u0430 \u0441\u0432\u043e\u044f \u0430\u043d\u043d\u043e\u0442\u0430\u0446\u0438\u044f \u043f\u043e \u0430\u043d\u0430\u043b\u043e\u0433\u0438\u0438 \u0441 \u0430\u043d\u043d\u043e\u0442\u0430\u0446\u0438\u0435\u0439 @WithFlinkCluster: <\/p>\n<pre><code class=\"java\">@Retention(RUNTIME) @ExtendWith({KafkaContainerExtension.class}) @Inherited public @interface WithKafkaContainer { }<\/code><\/pre>\n<p>\u042d\u0442\u0443 \u0430\u043d\u043d\u043e\u0442\u0430\u0446\u0438\u044e \u043c\u043e\u0436\u043d\u043e \u0432\u0435\u0448\u0430\u0442\u044c \u043d\u0430 \u043b\u044e\u0431\u043e\u0439 \u0442\u0435\u0441\u0442\u043e\u0432\u044b\u0439 \u043a\u043b\u0430\u0441\u0441, \u043a\u043e\u0442\u043e\u0440\u043e\u043c\u0443 \u043d\u0443\u0436\u043d\u0430 \u0438\u043d\u0442\u0435\u0433\u0440\u0430\u0446\u0438\u044f \u0441 Kafka.<\/p>\n<p><a class=\"anchor\" name=\"3\" id=\"3\"><\/a><\/p>\n<h2>\u0410\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u0438 \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u0434\u043b\u044f Kafka<\/h2>\n<p>\u0412\u043e \u0432\u0441\u0435\u0445 \u0442\u0435\u0441\u0442\u0430\u0445 \u0431\u044b\u043b\u043e \u0431\u044b \u0443\u0434\u043e\u0431\u043d\u043e \u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c\u0441\u044f \u0441\u0432\u043e\u0438\u043c\u0438 \u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u044f\u043c\u0438, \u0444\u0430\u0441\u0430\u0434\u0430\u043c\u0438 \u0438\u043b\u0438 dto \u0434\u043b\u044f \u0432\u0437\u0430\u0438\u043c\u043e\u0434\u0435\u0439\u0441\u0442\u0432\u0438\u044f \u0441 Kafka.<\/p>\n<p><a class=\"anchor\" name=\"4\" id=\"4\"><\/a><\/p>\n<h4>KafkaTestConsumer<\/h4>\n<p>\u0412\u043e-\u043f\u0435\u0440\u0432\u044b\u0445, \u043d\u0443\u0436\u043d\u043e \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0438\u0442\u044c \u0442\u0435\u0441\u0442\u043e\u0432\u044b\u0439 Consumer, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u0431\u0443\u0434\u0435\u043c \u0441\u043e\u0437\u0434\u0430\u0432\u0430\u0442\u044c \u0432 \u043a\u0430\u0436\u0434\u043e\u043c \u0442\u0435\u0441\u0442\u0435 \u043e\u0442\u0434\u0435\u043b\u044c\u043d\u043e, \u0447\u0442\u043e\u0431\u044b \u043e\u043d \u043f\u043e\u0434\u043a\u043b\u044e\u0447\u0430\u043b\u0441\u044f \u043a \u0442\u043e\u043f\u0438\u043a\u0430\u043c \u0432 \u0440\u0430\u043c\u043a\u0430\u0445 \u043d\u043e\u0432\u043e\u0439 \u043a\u043e\u043d\u0441\u044c\u044e\u043c\u0435\u0440\u043d\u043e\u0439 \u0433\u0440\u0443\u043f\u043f\u044b. \u0415\u0449\u0451 \u043e\u043d \u0434\u043e\u043b\u0436\u0435\u043d \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c \u0434\u0435-\/\u0441\u0435\u0440\u0438\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u044e Jackson, \u0432\u0435\u0434\u044c \u043c\u044b \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0438\u043b\u0438 \u0434\u043b\u044f \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0439 \u0444\u043e\u0440\u043c\u0430\u0442 JSON. \u0412\u0430\u0436\u043d\u043e \u043f\u043e\u043c\u043d\u0438\u0442\u044c: Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440 \u043f\u043e\u0434\u043d\u0438\u043c\u0430\u0435\u0442\u0441\u044f \u0432 \u0435\u0434\u0438\u043d\u0441\u0442\u0432\u0435\u043d\u043d\u043e\u043c \u044d\u043a\u0437\u0435\u043c\u043f\u043b\u044f\u0440\u0435. \u041f\u043e\u044d\u0442\u043e\u043c\u0443, \u0435\u0441\u043b\u0438 \u043d\u0430\u043f\u0438\u0441\u0430\u0442\u044c \u0442\u0435\u0441\u0442\u044b \u043d\u0435\u043f\u0440\u0430\u0432\u0438\u043b\u044c\u043d\u043e, \u043e\u043d\u0438 \u043c\u043e\u0433\u0443\u0442 \u043a\u043e\u0441\u0432\u0435\u043d\u043d\u043e \u0432\u043b\u0438\u044f\u0442\u044c \u0434\u0440\u0443\u0433 \u043d\u0430 \u0434\u0440\u0443\u0433\u0430. <\/p>\n<p>\u041a\u043b\u0430\u0441\u0441 \u0434\u043b\u044f \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u044f \u0442\u0435\u0441\u0442\u043e\u0432\u043e\u0433\u043e Kafka consumer \u043c\u043e\u0436\u043d\u043e \u043f\u0440\u0435\u0434\u0441\u0442\u0430\u0432\u0438\u0442\u044c \u0432 \u0442\u0430\u043a\u043e\u043c \u0432\u0438\u0434\u0435:<\/p>\n<pre><code class=\"java\">public class KafkaTestConsumer implements AutoCloseable {    private final Consumer&lt;String, String> consumer;    private final List&lt;KafkaMessage> receivedMessages = new CopyOnWriteArrayList&lt;>();    private final ObjectMapper objectMapper = createObjectMapper();     public KafkaTestConsumer(String bootstrapServers, Set&lt;String> topics) {        this.consumer = new KafkaConsumer&lt;>(            Map.of(                ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers,                ConsumerConfig.GROUP_ID_CONFIG, UUID.randomUUID().toString(),                ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName(),                ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName(),                ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 1,                ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, \"earliest\"            )        );        consumer.subscribe(topics);    }     public &lt;T> List&lt;T> receiveAndGetAll(String topic, Class&lt;T> clazz) {        return receiveAndGetAll()                   .stream()                   .filter(kafkaMessage -> topic.equals(kafkaMessage.getTopic()))                   .map(kafkaMessage -> readValue(kafkaMessage, clazz))                   .collect(toList());    }     private List&lt;KafkaMessage> receiveAndGetAll() {        final var records = consumer.poll(Duration.ofSeconds(5));        for (ConsumerRecord&lt;String, String> record : records) {            receivedMessages.add(new KafkaMessage(record.key(), record.topic(), record.value()));        }        consumer.commitSync();        return receivedMessages;    }     @SneakyThrows    private &lt;T> T readValue(KafkaMessage kafkaMessage, Class&lt;T> clazz) {        return objectMapper.readValue(kafkaMessage.getValue(), clazz);    }     @Override    public void close() {        receivedMessages.clear();        consumer.close();    } }<\/code><\/pre>\n<p>\u041a\u043e\u0434 \u0434\u043e\u0441\u0442\u0430\u0442\u043e\u0447\u043d\u043e \u043f\u0440\u043e\u0441\u0442. \u0412 \u043a\u043e\u043d\u0441\u0442\u0440\u0443\u043a\u0442\u043e\u0440\u0435 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u044f\u0435\u043c \u043e\u0441\u043d\u043e\u0432\u043d\u044b\u0435 \u043f\u0430\u0440\u0430\u043c\u0435\u0442\u0440\u044b \u043f\u043e\u0434\u043a\u043b\u044e\u0447\u0435\u043d\u0438\u044f \u043a Kafka, \u0441\u043e\u0437\u0434\u0430\u0451\u0442\u0441\u044f \u0431\u0430\u0437\u043e\u0432\u044b\u0439 consumer, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043f\u043e\u0434\u043f\u0438\u0441\u044b\u0432\u0430\u0435\u0442\u0441\u044f \u043d\u0430 \u043f\u0435\u0440\u0435\u0434\u0430\u043d\u043d\u044b\u0435 \u0442\u043e\u043f\u0438\u043a\u0438. \u041f\u043e\u0442\u043e\u043c \u043f\u0440\u0435\u0434\u043e\u0441\u0442\u0430\u0432\u043b\u044f\u0435\u043c \u043e\u0441\u043d\u043e\u0432\u043d\u043e\u0439 \u043c\u0435\u0442\u043e\u0434 receiveAndGetAll(), \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043e\u0442\u0434\u0430\u0451\u0442 \u0432\u0441\u0435 \u0434\u0435\u0441\u0435\u0440\u0438\u0430\u043b\u0438\u0437\u043e\u0432\u0430\u043d\u043d\u044b\u0435 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f \u0432 \u043a\u043e\u043d\u043a\u0440\u0435\u0442\u043d\u044b\u0439 \u0442\u0438\u043f \u0432 \u0440\u0430\u0437\u0440\u0435\u0437\u0435 \u043a\u0430\u043a\u043e\u0433\u043e-\u0442\u043e \u0442\u043e\u043f\u0438\u043a\u0430. \u0412\u043d\u0443\u0442\u0440\u0438 \u044d\u0442\u043e\u0433\u043e \u043c\u0435\u0442\u043e\u0434\u0430 \u043f\u0440\u043e\u0438\u0441\u0445\u043e\u0434\u0438\u0442 \u0432\u044b\u0437\u043e\u0432 \u0431\u0430\u0437\u043e\u0432\u043e\u0433\u043e consumer.poll() \u0434\u043b\u044f \u043e\u0431\u0440\u0430\u0449\u0435\u043d\u0438\u044f \u043a Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440\u0443. \u0412\u043d\u0443\u0442\u0440\u0438 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u043c \u0441\u0432\u043e\u044e dto KafkaMessage, \u0447\u0442\u043e\u0431\u044b \u043d\u0435 \u0440\u0430\u0431\u043e\u0442\u0430\u0442\u044c \u0441 ConsumerRecord \u043d\u0430\u043f\u0440\u044f\u043c\u0443\u044e, \u0432\u0435\u0434\u044c \u0431\u043e\u043b\u044c\u0448\u0438\u043d\u0441\u0442\u0432\u043e \u0438\u043d\u0444\u043e\u0440\u043c\u0430\u0446\u0438\u0438 \u0432 \u043d\u0451\u043c \u043d\u0430\u043c \u043d\u0435 \u0442\u0440\u0435\u0431\u0443\u0435\u0442\u0441\u044f:<\/p>\n<pre><code class=\"java\">@Value public class KafkaMessage {    String key;    String topic;    String value; }<\/code><\/pre>\n<p><a class=\"anchor\" name=\"5\" id=\"5\"><\/a><\/p>\n<h4>TestKafkaFacade<\/h4>\n<p>\u041d\u0430\u043c \u043f\u043e\u043d\u0430\u0434\u043e\u0431\u0438\u0442\u0441\u044f \u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u044f-\u0444\u0430\u0441\u0430\u0434, \u043a\u043e\u0442\u043e\u0440\u0430\u044f \u043f\u043e\u0437\u0432\u043e\u043b\u0438\u0442 \u043e\u0442\u043f\u0440\u0430\u0432\u043b\u044f\u0442\u044c \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f \u0432 Kafka, \u0441\u043e\u0437\u0434\u0430\u0432\u0430\u0442\u044c \u043d\u043e\u0432\u044b\u0435 \u0442\u043e\u043f\u0438\u043a\u0438 \u0432 \u043a\u0430\u0436\u0434\u043e\u043c \u0442\u0435\u0441\u0442\u0435, \u0430 \u0435\u0449\u0451 \u0441\u043e\u0437\u0434\u0430\u0432\u0430\u0442\u044c KafkaTestConsumer. \u042d\u0442\u043e \u0431\u0443\u0434\u0435\u0442 \u043f\u0440\u043e\u0438\u0441\u0445\u043e\u0434\u0438\u0442\u044c \u0432 \u043a\u043e\u043d\u0442\u0435\u043a\u0441\u0442\u0435 Spring-\u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u044f, \u043f\u043e\u044d\u0442\u043e\u043c\u0443 \u0443\u0434\u043e\u0431\u043d\u043e \u0441\u043e\u0437\u0434\u0430\u0442\u044c \u043e\u0431\u0449\u0438\u0439 Spring-\u043a\u043e\u043c\u043f\u043e\u043d\u0435\u043d\u0442 TestKafkaFacade:<\/p>\n<pre><code class=\"java\">@TestComponent @SuppressWarnings(\"PMD.TestClassWithoutTestCases\") public class TestKafkaFacade {    @Autowired    private KafkaTemplate&lt;String, String> kafkaTemplate;    @Autowired    private KafkaAdmin kafkaAdmin;     private final ObjectMapper objectMapper = createObjectMapper();     public void createTopicsIfNeeded(String... names) {        final var topics = kafkaAdmin.describeTopics();        if (!topics.keySet().containsAll(Stream.of(names).collect(toSet()))) {            kafkaAdmin.createOrModifyTopics(                Stream.of(names)                    .map(n -> new NewTopic(n, 1, (short) 1))                    .toArray(NewTopic[]::new)            );        }    }     @SneakyThrows    public void sendMessage(String topic, Object message) {        kafkaTemplate            .send(topic, objectMapper.writeValueAsString(message))            .get(5, TimeUnit.SECONDS);    }     public KafkaTestConsumer createKafkaConsumer(Set&lt;String> topics) {        return new KafkaTestConsumer(System.getProperty(\"kafka.bootstrap-servers\"), topics);    } }\u0412\u043e-\u043f\u0435\u0440\u0432\u044b\u0445, \u044d\u0442\u043e\u0442 \u0444\u0430\u0441\u0430\u0434 \u0438\u043d\u0436\u0435\u043a\u0442\u0438\u0442 Spring-\u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438 \u0434\u043b\u044f KafkaTemplate (\u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u044f \u043d\u0430\u0434 \u0431\u0430\u0437\u043e\u0432\u044b\u043c Kafka Producer), \u0430 \u0442\u0430\u043a\u0436\u0435 KafkaAdmin (\u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u044f \u043d\u0430\u0434 \u0431\u0430\u0437\u043e\u0432\u044b\u043c AdminClient). \u042d\u0442\u0438 \u0431\u0438\u043d\u044b \u0431\u0443\u0434\u0443\u0442 \u0441\u0443\u0449\u0435\u0441\u0442\u0432\u043e\u0432\u0430\u0442\u044c \u0432 \u043a\u043e\u043d\u0442\u0435\u043a\u0441\u0442\u0435 \u0431\u043b\u0430\u0433\u043e\u0434\u0430\u0440\u044f \u0443\u0441\u0442\u0430\u043d\u043e\u0432\u043b\u0435\u043d\u043d\u043e\u043c\u0443 \u0441\u0432\u043e\u0439\u0441\u0442\u0432\u0443 spring.kafka.bootstrap-servers \u0432 \u043d\u0430\u0448\u0435\u043c KafkaContainerExtension. \u0423 \u043d\u0430\u0441 \u0435\u0441\u0442\u044c \u043c\u0435\u0442\u043e\u0434\u044b \u0434\u043b\u044f \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u044f \u043d\u043e\u0432\u044b\u0445 \u0442\u043e\u043f\u0438\u043a\u043e\u0432, \u043e\u0442\u043f\u0440\u0430\u0432\u043a\u0438 \u0432 \u0442\u043e\u043f\u0438\u043a \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f \u0441 \u0443\u0447\u0451\u0442\u043e\u043c \u0441\u0435\u0440\u0438\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438 \u0438 \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u044f \u043d\u043e\u0432\u043e\u0433\u043e KafkaTestConsumer. Bootstrap-servers \u043f\u043e\u043b\u0443\u0447\u0430\u0435\u043c \u043f\u0440\u044f\u043c\u043e \u0438\u0437 env, \u0432 \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043c\u044b \u0437\u0430\u043f\u0438\u0441\u044b\u0432\u0430\u0435\u043c bootstrap-servers \u0432 \u0440\u0430\u043c\u043a\u0430\u0445 KafkaContainerExtension.   KafkaTopicCreatorConfig \u041f\u0440\u0435\u0436\u0434\u0435 \u0447\u0435\u043c \u043f\u0435\u0440\u0435\u0445\u043e\u0434\u0438\u0442\u044c \u043a \u0442\u0435\u0441\u0442\u0443, \u0443\u0434\u043e\u0431\u043d\u043e \u0432\u043e\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c\u0441\u044f \u0435\u0449\u0451 \u043e\u0434\u043d\u0438\u043c \u0442\u0435\u0441\u0442\u043e\u0432\u044b\u043c \u043a\u043e\u043c\u043f\u043e\u043d\u0435\u043d\u0442\u043e\u043c \u0434\u043b\u044f \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u044f \u0432\u0441\u0435\u0445 \u0442\u043e\u043f\u0438\u043a\u043e\u0432, \u0443\u043a\u0430\u0437\u0430\u043d\u043d\u044b\u0445 \u0432 application.yml: <\/code><\/pre>\n<p>\u0412\u043e-\u043f\u0435\u0440\u0432\u044b\u0445, \u044d\u0442\u043e\u0442 \u0444\u0430\u0441\u0430\u0434 \u0438\u043d\u0436\u0435\u043a\u0442\u0438\u0442 Spring-\u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438 \u0434\u043b\u044f KafkaTemplate (\u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u044f \u043d\u0430\u0434 \u0431\u0430\u0437\u043e\u0432\u044b\u043c Kafka Producer), \u0430 \u0442\u0430\u043a\u0436\u0435 KafkaAdmin (\u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u044f \u043d\u0430\u0434 \u0431\u0430\u0437\u043e\u0432\u044b\u043c AdminClient). \u042d\u0442\u0438 \u0431\u0438\u043d\u044b \u0431\u0443\u0434\u0443\u0442 \u0441\u0443\u0449\u0435\u0441\u0442\u0432\u043e\u0432\u0430\u0442\u044c \u0432 \u043a\u043e\u043d\u0442\u0435\u043a\u0441\u0442\u0435 \u0431\u043b\u0430\u0433\u043e\u0434\u0430\u0440\u044f \u0443\u0441\u0442\u0430\u043d\u043e\u0432\u043b\u0435\u043d\u043d\u043e\u043c\u0443 \u0441\u0432\u043e\u0439\u0441\u0442\u0432\u0443 spring.kafka.bootstrap-servers \u0432 \u043d\u0430\u0448\u0435\u043c KafkaContainerExtension.<\/p>\n<p>\u0423 \u043d\u0430\u0441 \u0435\u0441\u0442\u044c \u043c\u0435\u0442\u043e\u0434\u044b \u0434\u043b\u044f \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u044f \u043d\u043e\u0432\u044b\u0445 \u0442\u043e\u043f\u0438\u043a\u043e\u0432, \u043e\u0442\u043f\u0440\u0430\u0432\u043a\u0438 \u0432 \u0442\u043e\u043f\u0438\u043a \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f \u0441 \u0443\u0447\u0451\u0442\u043e\u043c \u0441\u0435\u0440\u0438\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438 \u0438 \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u044f \u043d\u043e\u0432\u043e\u0433\u043e KafkaTestConsumer. Bootstrap-servers \u043f\u043e\u043b\u0443\u0447\u0430\u0435\u043c \u043f\u0440\u044f\u043c\u043e \u0438\u0437 env, \u0432 \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043c\u044b \u0437\u0430\u043f\u0438\u0441\u044b\u0432\u0430\u0435\u043c bootstrap-servers \u0432 \u0440\u0430\u043c\u043a\u0430\u0445 KafkaContainerExtension.<\/p>\n<p><a class=\"anchor\" name=\"6\" id=\"6\"><\/a><\/p>\n<h4>KafkaTopicCreatorConfig<\/h4>\n<p>\u041f\u0440\u0435\u0436\u0434\u0435 \u0447\u0435\u043c \u043f\u0435\u0440\u0435\u0445\u043e\u0434\u0438\u0442\u044c \u043a \u0442\u0435\u0441\u0442\u0443, \u0443\u0434\u043e\u0431\u043d\u043e \u0432\u043e\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c\u0441\u044f \u0435\u0449\u0451 \u043e\u0434\u043d\u0438\u043c \u0442\u0435\u0441\u0442\u043e\u0432\u044b\u043c \u043a\u043e\u043c\u043f\u043e\u043d\u0435\u043d\u0442\u043e\u043c \u0434\u043b\u044f \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u044f \u0432\u0441\u0435\u0445 \u0442\u043e\u043f\u0438\u043a\u043e\u0432, \u0443\u043a\u0430\u0437\u0430\u043d\u043d\u044b\u0445 \u0432 application.yml:<\/p>\n<pre><code class=\"java\">@TestConfiguration public class KafkaTopicCreatorConfig {    @Autowired    private KafkaProperties kafkaProperties;     @Bean    public KafkaAdmin.NewTopics newTopics() {        return new KafkaAdmin.NewTopics(            new NewTopic(kafkaProperties.getTopics().getClickTopic(), 1, (short) 1)        );    } }<\/code><\/pre>\n<p>\u042d\u0442\u043e\u0442 \u043a\u043e\u043c\u043f\u043e\u043d\u0435\u043d\u0442 \u0441\u043e\u0437\u0434\u0430\u0451\u0442 \u0431\u0438\u043d KafkaAdmin.NewTopics, \u0430 Spring \u0430\u0432\u0442\u043e\u043c\u0430\u0442\u0438\u0447\u0435\u0441\u043a\u0438 \u0441\u043e\u0437\u0434\u0430\u0451\u0442 \u0432\u0441\u0435 \u0443\u043a\u0430\u0437\u0430\u043d\u043d\u044b\u0435 \u0432 \u043d\u0451\u043c \u0442\u043e\u043f\u0438\u043a\u0438 \u2014 \u0442\u043e \u0435\u0441\u0442\u044c \u0442\u043e\u043f\u0438\u043a\u0438 \u0438\u0437 KafkaProperties (application.yml). \u0412 \u043d\u0430\u0448\u0435\u043c \u0441\u043b\u0443\u0447\u0430\u0435 \u0441\u0435\u0439\u0447\u0430\u0441 \u044d\u0442\u043e \u0442\u043e\u043b\u044c\u043a\u043e \u043e\u0434\u0438\u043d \u0432\u0445\u043e\u0434\u043d\u043e\u0439 \u0442\u043e\u043f\u0438\u043a click-topic. \u0412\u044b\u0445\u043e\u0434\u043d\u044b\u0435 \u043c\u044b \u0431\u0443\u0434\u0435\u043c \u0441\u043e\u0437\u0434\u0430\u0432\u0430\u0442\u044c \u0432 \u043a\u0430\u0436\u0434\u043e\u043c \u0442\u0435\u0441\u0442\u0435 \u043e\u0442\u0434\u0435\u043b\u044c\u043d\u043e, \u043f\u043e\u0442\u043e\u043c\u0443 \u0447\u0442\u043e \u043e\u043d\u0438 \u043c\u043e\u0433\u0443\u0442 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u044f\u0442\u044c\u0441\u044f \u0434\u0438\u043d\u0430\u043c\u0438\u0447\u0435\u0441\u043a\u0438 \u0432 \u043d\u0430\u0448\u0435\u043c \u0431\u0438\u0437\u043d\u0435\u0441-\u043a\u0435\u0439\u0441\u0435.<\/p>\n<p><a class=\"anchor\" name=\"7\" id=\"7\"><\/a><\/p>\n<h2>E2E-\u0442\u0435\u0441\u0442 \u043d\u0430 Flink Job<\/h2>\n<p>\u0422\u0435\u043f\u0435\u0440\u044c \u0432\u0441\u0451 \u0433\u043e\u0442\u043e\u0432\u043e \u0434\u043b\u044f \u043d\u0430\u043f\u0438\u0441\u0430\u043d\u0438\u044f \u043f\u0435\u0440\u0432\u043e\u0433\u043e E2E-\u0442\u0435\u0441\u0442\u0430 \u043d\u0430 \u043d\u0430\u0448\u0443 \u0434\u0436\u043e\u0431\u0443. \u0414\u043b\u044f \u044d\u0442\u043e\u0433\u043e \u0441\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u043e\u0431\u0449\u0443\u044e \u0434\u043b\u044f E2E-\u0442\u0435\u0441\u0442\u043e\u0432 \u0430\u043d\u043d\u043e\u0442\u0430\u0446\u0438\u044e:<\/p>\n<pre><code class=\"java\">@Retention(RUNTIME) @SpringBootTest(webEnvironment = NONE) @ActiveProfiles({\"test\"}) @Import({         KafkaTopicCreatorConfig.class,         TestKafkaFacade.class }) @WithFlinkCluster @WithKafkaContainer public @interface E2ETest { }<\/code><\/pre>\n<p>\u041e\u043d\u0430 \u0441\u043e\u0432\u043c\u0435\u0449\u0430\u0435\u0442 \u0432 \u0441\u0435\u0431\u0435 \u0432\u0441\u0435 \u0438\u043d\u0442\u0435\u0433\u0440\u0430\u0446\u0438\u0438 \u043f\u0440\u043e\u0435\u043a\u0442\u0430 @WithFlinkCluster \u0438 @WithKafkaContainer. \u0410 \u0435\u0449\u0451 \u043f\u043e\u0434\u043d\u0438\u043c\u0430\u0435\u0442 \u0432\u0435\u0441\u044c Spring-\u043a\u043e\u043d\u0442\u0435\u043a\u0441\u0442, \u0437\u0430\u0445\u0432\u0430\u0442\u044b\u0432\u0430\u044f \u0442\u0435\u0441\u0442\u043e\u0432\u044b\u0435 \u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u0438, \u043a\u043e\u0442\u043e\u0440\u044b\u0435 \u043c\u044b \u043e\u043f\u0438\u0441\u0430\u043b\u0438 \u0432\u044b\u0448\u0435: KafkaTopicCreatorConfig \u0438 TestKafkaFacade.<\/p>\n<p>\u041d\u0430\u043f\u043e\u043c\u043d\u044e, \u0447\u0442\u043e \u0432 \u043d\u0430\u0448\u0435\u043c \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0438 \u0442\u043e\u0447\u043a\u0430 \u0432\u0445\u043e\u0434\u0430 \u2014 \u044d\u0442\u043e AppListener: <\/p>\n<pre><code class=\"java\">@Component @RequiredArgsConstructor @ConditionalOnProperty(\"flink.submit-jobs-on-app-start\") public class AppListener {    private final JobStarter jobStarter;     @EventListener(ApplicationStartedEvent.class)    @SneakyThrows    public void onApplicationStart() {        jobStarter.startJobs();    } }<\/code><\/pre>\n<p>\u041e\u043d \u0437\u0430\u043f\u0443\u0441\u043a\u0430\u0435\u0442\u0441\u044f \u043f\u043e \u0443\u0441\u043b\u043e\u0432\u0438\u044e flink.submit-jobs-on-app-start = true. \u041d\u043e \u0432 \u0442\u0435\u0441\u0442\u0430\u0445 \u043c\u044b \u0432\u044b\u0441\u0442\u0430\u0432\u0438\u043c \u0435\u0433\u043e \u0432 false (\u0432 application-test.yml), \u0447\u0442\u043e\u0431\u044b \u0441\u043e\u0437\u0434\u0430\u0432\u0430\u0442\u044c \u043d\u0435\u043e\u0431\u0445\u043e\u0434\u0438\u043c\u044b\u0435 \u043f\u0440\u0435\u0434\u0443\u0441\u043b\u043e\u0432\u0438\u044f \u043f\u0435\u0440\u0435\u0434 \u043d\u0435\u043f\u043e\u0441\u0440\u0435\u0434\u0441\u0442\u0432\u0435\u043d\u043d\u044b\u043c \u0437\u0430\u043f\u0443\u0441\u043a\u043e\u043c \u0442\u0435\u0441\u0442\u0430.<\/p>\n<p>\u041d\u0430\u043f\u043e\u043c\u043d\u044e, \u0447\u0442\u043e \u043d\u0430\u0448\u0430 \u0434\u0436\u043e\u0431\u0430 \u0437\u0430\u043d\u0438\u043c\u0430\u0435\u0442\u0441\u044f \u0444\u0438\u043b\u044c\u0442\u0440\u0430\u0446\u0438\u0435\u0439 \u0432\u0445\u043e\u0434\u043d\u043e\u0433\u043e \u043f\u043e\u0442\u043e\u043a\u0430 ClickMessage, \u043f\u0440\u043e\u043f\u0443\u0441\u043a\u0430\u044f \u0441\u043e\u0431\u044b\u0442\u0438\u044f \u0442\u043e\u043b\u044c\u043a\u043e \u0441 \u0442\u0438\u043f\u043e\u043c \u043f\u043b\u0430\u0442\u0444\u043e\u0440\u043c\u044b WEB \u0438 APP. \u0414\u0430\u043b\u044c\u0448\u0435 \u0441\u043e\u0431\u044b\u0442\u0438\u044f APP \u043f\u0440\u043e\u0445\u043e\u0434\u044f\u0442 \u0434\u0435\u0434\u0443\u043f\u043b\u0438\u043a\u0430\u0446\u0438\u044e. \u041f\u043e\u0442\u043e\u043c \u0441\u043e\u0431\u044b\u0442\u0438\u044f APP \u0438 WEB \u0437\u0430\u043f\u0438\u0441\u044b\u0432\u0430\u044e\u0442\u0441\u044f \u0432 \u0432\u044b\u0445\u043e\u0434\u043d\u043e\u0439 \u0442\u043e\u043f\u0438\u043a Kafka, \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u044f\u044e\u0449\u0438\u0439\u0441\u044f \u0434\u0438\u043d\u0430\u043c\u0438\u0447\u0435\u0441\u043a\u0438 \u0438\u0437 \u043f\u043e\u043b\u044f ClickMessage.productTopic. <\/p>\n<p>\u041f\u0430\u0439\u043f\u043b\u0430\u0439\u043d \u0432\u044b\u0433\u043b\u044f\u0434\u0438\u0442 \u0442\u0430\u043a:<\/p>\n<figure class=\"full-width\"><img loading=\"lazy\" decoding=\"async\" src=\"https:\/\/habrastorage.org\/r\/w1560\/getpro\/habr\/upload_files\/f94\/65b\/259\/f9465b25942f778fc78a3a6825fbf864.png\" width=\"3720\" height=\"462\" data-src=\"https:\/\/habrastorage.org\/getpro\/habr\/upload_files\/f94\/65b\/259\/f9465b25942f778fc78a3a6825fbf864.png\"\/><\/figure>\n<p>\u0412 \u0438\u0442\u043e\u0433\u0435 E2E-\u0442\u0435\u0441\u0442 \u043d\u0430 \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0443 ClickMessage-\u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f \u043c\u043e\u0436\u0435\u0442 \u0432\u044b\u0433\u043b\u044f\u0434\u0435\u0442\u044c \u0442\u0430\u043a\u0438\u043c \u043e\u0431\u0440\u0430\u0437\u043e\u043c:<\/p>\n<pre><code class=\"java\">@E2ETest @SuppressWarnings(\"PMD.DataflowAnomalyAnalysis\") public class JobE2ETest {    @Autowired    private JobStarter jobStarter;    @Autowired    private TestKafkaFacade kafka;    @Autowired    private KafkaProperties kafkaProperties;     @Test    @SneakyThrows    void shouldProcessClickMessageSourceToProductSink() {        final var productTopic = \"product_topic_1\";        kafka.createTopicsIfNeeded(productTopic);        final var clickMessage = aClickMessage().withProductTopic(productTopic).withPlatform(APP).build();        kafka.sendMessage(kafkaProperties.getTopics().getClickTopic(), clickMessage);        kafka.sendMessage(kafkaProperties.getTopics().getClickTopic(), clickMessage);         final var jobClient = jobStarter.startJobs();         @Cleanup final var kafkaConsumer =            kafka.createKafkaConsumer(Set.of(productTopic));        await().atMost(ofSeconds(5))            .until(() -> kafkaConsumer.receiveAndGetAll(productTopic, ProductMessage.class),                productMessages -> productMessages.size() == 1                                       &amp;&amp; productMessages.get(0).getUserId().equals(clickMessage.getUserId())            );        jobClient.cancel().get(5, TimeUnit.SECONDS);    } }<\/code><\/pre>\n<p>\u041c\u044b \u043f\u043e\u0432\u0435\u0441\u0438\u043b\u0438 \u0435\u0434\u0438\u043d\u0441\u0442\u0432\u0435\u043d\u043d\u0443\u044e \u0430\u043d\u043d\u043e\u0442\u0430\u0446\u0438\u044e @E2ETest, \u0431\u043b\u0430\u0433\u043e\u0434\u0430\u0440\u044f \u043a\u043e\u0442\u043e\u0440\u043e\u0439 \u043f\u0440\u043e\u0438\u0437\u043e\u0439\u0434\u0451\u0442 \u0437\u0430\u043f\u0443\u0441\u043a Flink-\u043c\u0438\u043d\u0438-\u043a\u043b\u0430\u0441\u0442\u0435\u0440\u0430, Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440\u0430 \u0438 \u0432\u0441\u0435\u0433\u043e Spring-\u043a\u043e\u043d\u0442\u0435\u043a\u0441\u0442\u0430 \u043d\u0430\u0448\u0435\u0433\u043e \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u044f. \u0422\u043e\u0447\u043a\u0443 \u0432\u0445\u043e\u0434\u0430 JobStarter \u0431\u0443\u0434\u0435\u043c \u0437\u0430\u043f\u0443\u0441\u043a\u0430\u0442\u044c \u0441\u0430\u043c\u0438 \u0432 \u043e\u0431\u0445\u043e\u0434 AppListener \u0443\u0436\u0435 \u043f\u043e\u0441\u043b\u0435 \u043d\u0430\u0441\u0442\u0440\u043e\u0439\u043a\u0438 \u0442\u0435\u0441\u0442\u0430, \u0431\u043b\u0430\u0433\u043e\u0434\u0430\u0440\u044f \u043d\u0430\u0441\u0442\u0440\u043e\u0439\u043a\u0435 \u0432 application-test.yml: flink.submit-jobs-on-app-start = false.<\/p>\n<p>\u0410 \u0432\u043e\u0442 \u0447\u0442\u043e \u043f\u0440\u043e\u0438\u0441\u0445\u043e\u0434\u0438\u0442 \u0432 \u0441\u0430\u043c\u043e\u043c \u0442\u0435\u0441\u0442\u0435:<\/p>\n<ol>\n<li>\n<p>\u0421\u043e\u0437\u0434\u0430\u0451\u043c \u043d\u043e\u0432\u044b\u0439 productTopic. \u041c\u044b \u043e\u0436\u0438\u0434\u0430\u0435\u043c, \u0447\u0442\u043e \u0438\u043c\u0435\u043d\u043d\u043e \u0432 \u043d\u0435\u0433\u043e \u043f\u043e\u043f\u0430\u0434\u0451\u0442 \u0432\u0445\u043e\u0434\u043d\u043e\u0435 ClickMessage-\u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0435 \u043f\u043e\u0441\u043b\u0435 \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438.<\/p>\n<\/li>\n<li>\n<p>\u0421\u043e\u0437\u0434\u0430\u0451\u043c \u0441\u0430\u043c\u043e ClickMessage-\u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0435. \u041e\u0431\u044f\u0437\u0430\u0442\u0435\u043b\u044c\u043d\u043e \u043f\u0435\u0440\u0435\u0434\u0430\u0451\u043c \u0432 \u043d\u0435\u0433\u043e productTopic.<\/p>\n<\/li>\n<li>\n<p>\u041e\u0442\u043f\u0440\u0430\u0432\u043b\u044f\u0435\u043c \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0435 \u0432\u043e \u0432\u0445\u043e\u0434\u043d\u043e\u0439 \u0442\u043e\u043f\u0438\u043a click-topic \u0434\u0432\u0430\u0436\u0434\u044b (\u043e\u0436\u0438\u0434\u0430\u0435\u043c      \u0434\u0435\u0434\u0443\u043f\u043b\u0438\u043a\u0430\u0446\u0438\u044e).<\/p>\n<\/li>\n<\/ol>\n<p>\u041d\u0430 \u044d\u0442\u043e\u043c \u043f\u043e\u0434\u0433\u043e\u0442\u043e\u0432\u043a\u0430 \u0442\u0435\u0441\u0442\u0430 \u0437\u0430\u0432\u0435\u0440\u0448\u0435\u043d\u0430. \u0423 \u043d\u0430\u0441 \u0435\u0441\u0442\u044c:<\/p>\n<ul>\n<li>\n<p>\u0432\u0441\u0435 \u0442\u043e\u043f\u0438\u043a\u0438<\/p>\n<\/li>\n<li>\n<p>\u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0435 \u0432\u043e \u0432\u0445\u043e\u0434\u043d\u043e\u043c \u0442\u043e\u043f\u0438\u043a\u0435<\/p>\n<\/li>\n<li>\n<p>\u043f\u043e\u0434\u043d\u044f\u0442\u043e\u0435 Spring-\u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0435<\/p>\n<\/li>\n<\/ul>\n<p>\u041f\u043e\u0440\u0430 \u0441\u0442\u0430\u0440\u0442\u043e\u0432\u0430\u0442\u044c \u0437\u0430\u0434\u0430\u0447\u0443! \u041d\u0430\u0448\u0430 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u044f jobStarter.startJobs(); \u0437\u0430\u043f\u0443\u0441\u043a\u0430\u0435\u0442 \u0432\u0441\u0435 \u043d\u0430\u0439\u0434\u0435\u043d\u043d\u044b\u0435 \u0432 \u043a\u043e\u043d\u0442\u0435\u043a\u0441\u0442\u0435 Spring \u0437\u0430\u0434\u0430\u0447\u0438 (\u043d\u0430\u0441\u043b\u0435\u0434\u0443\u044e\u0449\u0438\u0435\u0441\u044f \u043e\u0442 \u043d\u0430\u0448\u0435\u0439 \u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u0438 FlinkJob) \u0430\u0441\u0438\u043d\u0445\u0440\u043e\u043d\u043d\u043e \u0447\u0435\u0440\u0435\u0437 environment.executeAsync(). \u0414\u0430\u043b\u044c\u0448\u0435 \u043c\u044b \u043c\u043e\u0436\u0435\u043c \u0434\u043e\u0436\u0438\u0434\u0430\u0442\u044c\u0441\u044f \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f \u0432 \u0432\u044b\u0445\u043e\u0434\u043d\u043e\u043c \u0442\u043e\u043f\u0438\u043a\u0435. <\/p>\n<p>\u0412\u0430\u0436\u043d\u044b\u0439 \u043c\u043e\u043c\u0435\u043d\u0442 \u2014 \u0438\u043c\u0435\u043d\u043d\u043e \u0430\u0441\u0438\u043d\u0445\u0440\u043e\u043d\u043d\u044b\u0439 \u0437\u0430\u043f\u0443\u0441\u043a \u0437\u0430\u0434\u0430\u043d\u0438\u044f, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u0432\u043e\u0437\u0432\u0440\u0430\u0449\u0430\u0435\u0442 \u043e\u0431\u044a\u0435\u043a\u0442 \u0443\u043f\u0440\u0430\u0432\u043b\u0435\u043d\u0438\u044f JobClient. \u041d\u0430\u043f\u0440\u0438\u043c\u0435\u0440, \u0440\u0430\u043d\u044c\u0448\u0435 \u0432 \u0442\u0435\u0441\u0442\u0430\u0445 \u043c\u044b \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u043b\u0438 \u0442\u043e\u043b\u044c\u043a\u043e \u043c\u0435\u0442\u043e\u0434 execute(), \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u0431\u044b\u043b \u0441\u0438\u043d\u0445\u0440\u043e\u043d\u043d\u044b\u043c. \u0412 \u0441\u043b\u0443\u0447\u0430\u0435 \u0430\u0441\u0438\u043d\u0445\u0440\u043e\u043d\u043d\u043e\u0433\u043e \u0437\u0430\u043f\u0443\u0441\u043a\u0430 \u0447\u0435\u0440\u0435\u0437 JobClient \u043c\u044b \u043c\u043e\u0436\u0435\u043c \u043f\u043e\u043b\u0443\u0447\u0430\u0442\u044c \u0441\u0442\u0430\u0442\u0443\u0441 \u0437\u0430\u0434\u0430\u0447\u0438, \u0437\u0430\u0432\u0435\u0440\u0448\u0430\u0442\u044c \u0435\u0451 \u0440\u0443\u043a\u0430\u043c\u0438 \u0438 \u0442\u0430\u043a \u0434\u0430\u043b\u0435\u0435. \u042d\u0442\u043e \u0432\u0430\u0436\u043d\u043e, \u0432\u0435\u0434\u044c \u043d\u0430\u0448 Kafka-\u0438\u0441\u0442\u043e\u0447\u043d\u0438\u043a \u043f\u043e\u0442\u0435\u043d\u0446\u0438\u0430\u043b\u044c\u043d\u043e \u0431\u0435\u0441\u043a\u043e\u043d\u0435\u0447\u0435\u043d.<\/p>\n<p>\u041f\u043e\u044d\u0442\u043e\u043c\u0443 \u0434\u0430\u043b\u044c\u0448\u0435 \u043c\u044b \u0441\u043e\u0437\u0434\u0430\u0451\u043c \u043d\u0430\u0448 AutoCloseable Kafka Consumer \u0438 \u043f\u043e\u0434\u043f\u0438\u0441\u044b\u0432\u0430\u0435\u043c \u0435\u0433\u043e \u043d\u0430 \u0432\u044b\u0445\u043e\u0434\u043d\u043e\u0439 productTopic. \u041f\u043e\u0442\u043e\u043c \u043f\u0435\u0440\u0438\u043e\u0434\u0438\u0447\u0435\u0441\u043a\u0438 \u043f\u0440\u043e\u0432\u0435\u0440\u044f\u0435\u043c, \u043f\u043e\u044f\u0432\u0438\u043b\u043e\u0441\u044c \u043b\u0438 \u0442\u0430\u043c \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0435 ProductMessage \u0441 userId, \u0438\u0434\u0435\u043d\u0442\u0438\u0447\u043d\u044b\u043c \u0432\u0445\u043e\u0434\u043d\u043e\u043c\u0443 ClickMessage, \u0432 \u0442\u0435\u0447\u0435\u043d\u0438\u0435 \u043f\u044f\u0442\u0438 \u0441\u0435\u043a\u0443\u043d\u0434. \u0415\u0441\u043b\u0438 \u0437\u0430 \u043f\u044f\u0442\u044c \u0441\u0435\u043a\u0443\u043d\u0434 \u043d\u0435 \u0434\u043e\u0436\u0434\u0430\u043b\u0438\u0441\u044c \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f, \u0442\u0435\u0441\u0442 \u0437\u0430\u0432\u0435\u0440\u0448\u0438\u0442\u0441\u044f \u0441 \u043e\u0448\u0438\u0431\u043a\u043e\u0439. \u0414\u043b\u044f \u043f\u0440\u043e\u0433\u0440\u0430\u043c\u043c\u043d\u043e\u0439 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438 \u0442\u0430\u043a\u043e\u0439 \u043f\u0440\u043e\u0432\u0435\u0440\u043a\u0438 \u0432 \u0442\u0435\u0441\u0442\u0435 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u044e\u0442\u0441\u044f \u0431\u0438\u0431\u043b\u0438\u043e\u0442\u0435\u043a\u0438 <a href=\"https:\/\/github.com\/awaitility\/awaitility\">awaitility<\/a>. \u041f\u043e\u0441\u043b\u0435 \u043f\u0440\u043e\u0432\u0435\u0440\u043a\u0438 \u0437\u0430\u0432\u0435\u0440\u0448\u0430\u0435\u043c \u0430\u0441\u0438\u043d\u0445\u0440\u043e\u043d\u043d\u0443\u044e \u0434\u0436\u043e\u0431\u0443 \u0441\u0438\u043d\u0445\u0440\u043e\u043d\u043d\u044b\u043c \u0432\u044b\u0437\u043e\u0432\u043e\u043c: jobClient.cancel().<\/p>\n<p><strong>\u0412 \u044d\u0442\u043e\u043c \u0442\u0435\u0441\u0442\u0435 \u0435\u0441\u0442\u044c \u0431\u043e\u043b\u044c\u0448\u0430\u044f<\/strong> <strong>\u043f\u0440\u043e\u0431\u043b\u0435\u043c\u0430<\/strong>: \u0430 \u0447\u0442\u043e, \u0435\u0441\u043b\u0438 assert \u043d\u0435 \u0432\u044b\u043f\u043e\u043b\u043d\u0438\u0442\u0441\u044f \u0438\u043b\u0438 \u0434\u0436\u043e\u0431\u0430 \u0432\u043e \u0432\u0440\u0435\u043c\u044f \u0437\u0430\u043f\u0443\u0441\u043a\u0430 \u0443\u043f\u0430\u0434\u0451\u0442? \u0411\u0443\u0434\u0443\u0442 \u043b\u0438 \u043a\u0430\u043a\u0438\u0435-\u0442\u043e \u043f\u0440\u043e\u0431\u043b\u0435\u043c\u044b \u0441 \u0441\u0430\u043c\u0438\u043c\u0438 \u0442\u0435\u0441\u0442\u0430\u043c\u0438?<\/p>\n<p>\u041d\u0430 \u0441\u0430\u043c\u043e\u043c \u0434\u0435\u043b\u0435 \u043f\u0440\u043e\u0431\u043b\u0435\u043c\u044b \u0431\u0443\u0434\u0443\u0442. \u041d\u0430\u043f\u0440\u0438\u043c\u0435\u0440, \u043c\u0435\u0442\u043e\u0434 jobClient.cancel() \u0432\u043e\u043e\u0431\u0449\u0435 \u043d\u0435 \u0432\u044b\u043f\u043e\u043b\u043d\u0438\u0442\u0441\u044f, \u0430 \u0434\u0436\u043e\u0431\u0430 \u0435\u0449\u0451 \u0434\u043e\u043b\u0433\u043e \u043c\u043e\u0436\u0435\u0442 \u0432\u0438\u0441\u0435\u0442\u044c \u0432 \u043c\u0438\u043d\u0438-\u043a\u043b\u0430\u0441\u0442\u0435\u0440\u0435. \u0412 \u044d\u0442\u043e \u0432\u0440\u0435\u043c\u044f \u043c\u0438\u043d\u0438-\u043a\u043b\u0430\u0441\u0442\u0435\u0440 \u043c\u043e\u0436\u0435\u0442 \u043d\u0430\u0447\u0430\u0442\u044c \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c\u0441\u044f \u0434\u043b\u044f \u0441\u043b\u0435\u0434\u0443\u044e\u0449\u0435\u0439 \u0434\u0436\u043e\u0431\u044b \u0441\u043b\u0435\u0434\u0443\u044e\u0449\u0435\u0433\u043e \u0442\u0435\u0441\u0442\u0430, \u0430 \u0435\u0433\u043e \u0440\u0435\u0441\u0443\u0440\u0441\u043e\u0432 \u0434\u043b\u044f \u044d\u0442\u043e\u0433\u043e \u043d\u0435 \u0445\u0432\u0430\u0442\u0438\u0442. \u041f\u043b\u044e\u0441 \u0432\u043e\u0437\u043c\u043e\u0436\u043d\u044b \u0440\u0430\u0437\u043d\u044b\u0435 \u0441\u0430\u0439\u0434-\u044d\u0444\u0444\u0435\u043a\u0442\u044b \u043c\u0435\u0436\u0434\u0443 \u0432\u044b\u043f\u043e\u043b\u043d\u0435\u043d\u0438\u0435\u043c \u0442\u0430\u043a\u0438\u0445 \u0442\u0435\u0441\u0442\u043e\u0432. \u041a\u0430\u043a \u0440\u0435\u0448\u0438\u0442\u044c \u044d\u0442\u0443 \u043f\u0440\u043e\u0431\u043b\u0435\u043c\u0443? \u041e\u0431 \u044d\u0442\u043e\u043c \u2014 \u0434\u0430\u043b\u044c\u0448\u0435.<\/p>\n<p><a class=\"anchor\" name=\"8\" id=\"8\"><\/a><\/p>\n<h2>\u0411\u0435\u0437\u043e\u043f\u0430\u0441\u043d\u043e\u0435 \u0437\u0430\u0432\u0435\u0440\u0448\u0435\u043d\u0438\u0435 E2E-\u0442\u0435\u0441\u0442\u043e\u0432<\/h2>\n<p>\u0427\u0442\u043e\u0431\u044b \u0440\u0435\u0448\u0438\u0442\u044c \u043f\u0440\u043e\u0431\u043b\u0435\u043c\u0443, \u043a\u043e\u0442\u043e\u0440\u0443\u044e \u044f \u043e\u043f\u0438\u0441\u0430\u043b \u0432\u044b\u0448\u0435, \u043f\u0435\u0440\u0435\u0434 \u043d\u0435\u0443\u0434\u0430\u0447\u043d\u044b\u043c \u0437\u0430\u0432\u0435\u0440\u0448\u0435\u043d\u0438\u0435\u043c \u0442\u0435\u0441\u0442\u0430 \u043b\u0443\u0447\u0448\u0435 \u0432\u0441\u0435\u0433\u0434\u0430 \u043f\u043e\u0434\u0447\u0438\u0449\u0430\u0442\u044c \u0440\u0435\u0441\u0443\u0440\u0441\u044b. \u0414\u043b\u044f \u044d\u0442\u043e\u0433\u043e \u043c\u043e\u0436\u043d\u043e \u0432 \u043b\u043e\u0431 \u043e\u0431\u043e\u0440\u0430\u0447\u0438\u0432\u0430\u0442\u044c \u0432\u0441\u0451 \u0432 \u0431\u043b\u043e\u043a try \u0438 finally. \u0412\u043e\u0442 \u043a\u0430\u043a \u044d\u0442\u043e \u0441\u0434\u0435\u043b\u0430\u0442\u044c:<\/p>\n<pre><code class=\"java\">try {     final var jobClient = jobStarter.startJobs();     \/\/ ...     await().atMost(ofSeconds(5)).until(...); } finally {     jobClient.cancel().get(5, TimeUnit.SECONDS); }<\/code><\/pre>\n<p>\u0412\u044b\u0433\u043b\u044f\u0434\u0438\u0442 \u043d\u0435 \u043e\u0447\u0435\u043d\u044c \u043a\u0440\u0430\u0441\u0438\u0432\u043e. \u0410 \u0435\u0449\u0451 \u0442\u0430\u043a\u043e\u0439 \u0431\u043b\u043e\u043a \u043f\u0440\u0438\u0434\u0451\u0442\u0441\u044f \u043f\u0438\u0441\u0430\u0442\u044c \u0432 \u043a\u0430\u0436\u0434\u043e\u043c \u0442\u0435\u0441\u0442\u0435, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u0442 \u0430\u0441\u0438\u043d\u0445\u0440\u043e\u043d\u043d\u044b\u0439 \u0437\u0430\u043f\u0443\u0441\u043a \u0437\u0430\u0434\u0430\u0447. \u041d\u0430\u043f\u0440\u0430\u0448\u0438\u0432\u0430\u0435\u0442\u0441\u044f \u0434\u0435\u043a\u043e\u0440\u0430\u0442\u043e\u0440 \u043d\u0430\u0434 \u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u0435\u0439 FlinkClient, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u0443\u043c\u0435\u0435\u0442 \u0437\u0430\u0432\u0435\u0440\u0448\u0430\u0442\u044c \u0437\u0430\u0434\u0430\u043d\u0438\u044f:<\/p>\n<pre><code class=\"java\">@RequiredArgsConstructor public class AutoCloseableJobClient implements JobClient, AutoCloseable {    private final JobClient original;     @Override    public JobID getJobID() {        return original.getJobID();    }     @Override    public CompletableFuture&lt;JobStatus> getJobStatus() {        return original.getJobStatus();    }     @Override    public CompletableFuture&lt;Void> cancel() {        return original.cancel();    }     @Override    public CompletableFuture&lt;String> stopWithSavepoint(boolean advanceToEndOfEventTime, @Nullable String savepointDirectory, SavepointFormatType formatType) {        return original.stopWithSavepoint(advanceToEndOfEventTime, savepointDirectory, formatType);    }     @Override    public CompletableFuture&lt;String> triggerSavepoint(@Nullable String savepointDirectory, SavepointFormatType formatType) {        return original.triggerSavepoint(savepointDirectory, formatType);    }     @Override    public CompletableFuture&lt;Map&lt;String, Object>> getAccumulators() {        return original.getAccumulators();    }     @Override    public CompletableFuture&lt;JobExecutionResult> getJobExecutionResult() {        return original.getJobExecutionResult();    }     @Override    public void close() throws Exception {        original.cancel().get(5, TimeUnit.SECONDS);    } }<\/code><\/pre>\n<p>\u0414\u0435\u043a\u043e\u0440\u0430\u0442\u043e\u0440 \u0434\u0435\u043b\u0435\u0433\u0438\u0440\u0443\u0435\u0442 \u0432\u044b\u043f\u043e\u043b\u043d\u0435\u043d\u0438\u0435 \u0432\u043d\u0443\u0442\u0440\u0435\u043d\u043d\u0435\u043c\u0443 \u043e\u0440\u0438\u0433\u0438\u043d\u0430\u043b\u044c\u043d\u043e\u043c\u0443 JobClient. \u041d\u043e \u0434\u043e\u043f\u043e\u043b\u043d\u0438\u0442\u0435\u043b\u044c\u043d\u043e \u043c\u044b \u0438\u043c\u043f\u043b\u0435\u043c\u0435\u043d\u0442\u0438\u0440\u0443\u0435\u043c \u0438\u043d\u0442\u0435\u0440\u0444\u0435\u0439\u0441 AutoCloseable, \u0447\u0442\u043e\u0431\u044b \u043f\u0435\u0440\u0435\u043d\u0435\u0441\u0442\u0438 \u043b\u043e\u0433\u0438\u043a\u0443 \u0437\u0430\u0432\u0435\u0440\u0448\u0435\u043d\u0438\u044f \u0434\u0436\u043e\u0431\u044b \u0432 \u043c\u0435\u0442\u043e\u0434 close(). \u0422\u0435\u043f\u0435\u0440\u044c \u043e\u0441\u0442\u0430\u043b\u043e\u0441\u044c \u0432\u0435\u0440\u043d\u0443\u0442\u044c \u044d\u0442\u043e\u0442 \u0434\u0435\u043a\u043e\u0440\u0430\u0442\u043e\u0440 \u0432 \u043d\u0430\u0448\u0435\u043c JobStarter:<\/p>\n<pre><code class=\"java\">@SneakyThrows public AutoCloseableJobClient startJobs() {    if (jobs.isEmpty()) {        log.info(\"No Jobs found for start\");        return null;    }    for (FlinkJob job : jobs) {        log.info(\"Register job '{}'\", job.getClass().getSimpleName());        job.registerJob(environment);    }    return new AutoCloseableJobClient(environment.executeAsync()); }<\/code><\/pre>\n<p>\u0422\u0435\u0441\u0442\u044b \u043f\u0440\u0438 \u044d\u0442\u043e\u043c \u0437\u043d\u0430\u0447\u0438\u0442\u0435\u043b\u044c\u043d\u043e \u0443\u043f\u0440\u043e\u0449\u0430\u044e\u0442\u0441\u044f:<\/p>\n<pre><code class=\"java\">@Test @SneakyThrows void shouldProcessClickMessageSourceToProductSink() {    final var productTopic = \"product_topic_1\";    kafka.createTopicsIfNeeded(productTopic);    final var clickMessage = aClickMessage().withProductTopic(productTopic).withPlatform(APP).build();    kafka.sendMessage(kafkaProperties.getTopics().getClickTopic(), clickMessage);    kafka.sendMessage(kafkaProperties.getTopics().getClickTopic(), clickMessage);     @Cleanup final var jobClient = jobStarter.startJobs();     @Cleanup final var kafkaConsumer =        kafka.createKafkaConsumer(Set.of(productTopic));    await().atMost(ofSeconds(5))        .until(() -> kafkaConsumer.receiveAndGetAll(productTopic, ProductMessage.class),            productMessages -> productMessages.size() == 1                                   &amp;&amp; productMessages.get(0).getUserId().equals(clickMessage.getUserId())        ); }<\/code><\/pre>\n<p>\u041f\u043e\u043b\u0443\u0447\u0430\u0435\u0442\u0441\u044f, \u043c\u044b \u0434\u043e\u0431\u0430\u0432\u0438\u043b\u0438 \u0430\u043d\u043d\u043e\u0442\u0430\u0446\u0438\u044e @Cleanup \u0434\u043b\u044f \u043e\u0431\u044a\u0435\u043a\u0442\u0430 JobClient \u0438 \u0443\u0431\u0440\u0430\u043b\u0438 \u043b\u0438\u0448\u043d\u044e\u044e \u0437\u0430\u0432\u0435\u0440\u0448\u0430\u044e\u0449\u0443\u044e \u0441\u0442\u0440\u043e\u043a\u0443 jobClient.cancel().get(5, TimeUnit.SECONDS).<\/p>\n<p><a class=\"anchor\" name=\"9\" id=\"9\"><\/a><\/p>\n<h2>RocksDB \u0432 E2E-\u0442\u0435\u0441\u0442\u0430\u0445<\/h2>\n<p>\u0427\u0442\u043e\u0431\u044b \u043e\u043a\u043e\u043d\u0447\u0430\u0442\u0435\u043b\u044c\u043d\u043e \u0443\u0431\u0435\u0434\u0438\u0442\u044c\u0441\u044f, \u0447\u0442\u043e \u043d\u0430\u0448\u0438 E2E-\u0442\u0435\u0441\u0442\u044b \u043f\u043e\u043a\u0440\u044b\u0432\u0430\u044e\u0442 \u043d\u0435\u043e\u0431\u0445\u043e\u0434\u0438\u043c\u0443\u044e \u0444\u0443\u043d\u043a\u0446\u0438\u043e\u043d\u0430\u043b\u044c\u043d\u043e\u0441\u0442\u044c, \u043d\u0443\u0436\u043d\u043e \u0432\u043e\u0441\u0441\u043e\u0437\u0434\u0430\u0442\u044c \u043e\u043a\u0440\u0443\u0436\u0435\u043d\u0438\u0435, \u0430\u043d\u0430\u043b\u043e\u0433\u0438\u0447\u043d\u043e\u0435 production. \u0412 \u043d\u0451\u043c \u0434\u043b\u044f \u0431\u043e\u043b\u044c\u0448\u0438\u0445 \u0441\u043e\u0441\u0442\u043e\u044f\u043d\u0438\u0439 \u0440\u0435\u043a\u043e\u043c\u0435\u043d\u0434\u0443\u0435\u0442\u0441\u044f \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c RocksDB \u0432 \u043a\u0430\u0447\u0435\u0441\u0442\u0432\u0435 StateBackend, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043f\u0440\u0438\u043c\u0435\u043d\u044f\u0435\u0442\u0441\u044f \u043d\u0430\u043f\u0440\u044f\u043c\u0443\u044e \u0432 \u0434\u0435\u0434\u0443\u043f\u043b\u0438\u043a\u0430\u0442\u043e\u0440\u0435. RocksDB \u0432 \u043a\u0430\u0447\u0435\u0441\u0442\u0432\u0435 \u0431\u044d\u043a\u0435\u043d\u0434\u0430 \u0441\u043e\u0441\u0442\u043e\u044f\u043d\u0438\u0439 \u044f \u043e\u043f\u0438\u0441\u044b\u0432\u0430\u043b \u0432 <a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/772898\/\">\u043f\u0435\u0440\u0432\u043e\u0439 \u0441\u0442\u0430\u0442\u044c\u0435<\/a> \u044d\u0442\u043e\u0433\u043e \u0446\u0438\u043a\u043b\u0430.<\/p>\n<p>\u0418\u0442\u0430\u043a, \u043c\u043e\u0436\u043d\u043e \u0437\u0430\u0441\u0442\u0430\u0432\u0438\u0442\u044c E2E-\u0442\u0435\u0441\u0442 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c \u043b\u043e\u043a\u0430\u043b\u044c\u043d\u043e \u043f\u043e\u0434\u043d\u044f\u0442\u044b\u0439 RocksDB. \u0421\u0434\u0435\u043b\u0430\u0442\u044c \u044d\u0442\u043e \u0434\u043e\u0441\u0442\u0430\u0442\u043e\u0447\u043d\u043e \u043f\u0440\u043e\u0441\u0442\u043e, \u0442\u0430\u043a \u043a\u0430\u043a Flink \u043f\u0440\u0435\u0434\u043e\u0441\u0442\u0430\u0432\u043b\u044f\u0435\u0442 RocksDB \u00ab\u0438\u0437 \u043a\u043e\u0440\u043e\u0431\u043a\u0438\u00bb. \u0412 \u0442\u0435\u0441\u0442\u0430\u0445 \u0434\u043e\u0441\u0442\u0430\u0442\u043e\u0447\u043d\u043e \u0443\u043a\u0430\u0437\u0430\u0442\u044c \u0435\u0433\u043e \u0432 \u043a\u0430\u0447\u0435\u0441\u0442\u0432\u0435 \u0441\u043e\u0441\u0442\u043e\u044f\u043d\u0438\u044f:<\/p>\n<pre><code class=\"java\">@TestConfiguration public class FlinkProductionConfig {     @Autowired    public void changeFlinkEnvironment(StreamExecutionEnvironment environment) {        final var backend = new EmbeddedRocksDBStateBackend(false);        environment.setStateBackend(backend);    } }<\/code><\/pre>\n<p>\u0422\u043e \u0435\u0441\u0442\u044c \u043d\u0430\u0448 \u043a\u043e\u043c\u043f\u043e\u043d\u0435\u043d\u0442 \u043f\u0435\u0440\u0435\u0445\u0432\u0430\u0442\u044b\u0432\u0430\u0435\u0442 \u043d\u0430\u0441\u0442\u0440\u043e\u0439\u043a\u0443 StreamExecutionEnvironment \u0438 \u0443\u043a\u0430\u0437\u044b\u0432\u0430\u0435\u0442 \u043d\u0430\u043f\u0440\u044f\u043c\u0443\u044e \u0441\u043e\u0441\u0442\u043e\u044f\u043d\u0438\u0435 EmbeddedRocksDBStateBackend. \u0427\u0442\u043e\u0431\u044b \u0435\u0433\u043e \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c, \u043d\u0443\u0436\u043d\u043e \u0434\u043e\u0431\u0430\u0432\u0438\u0442\u044c \u0437\u0430\u0432\u0438\u0441\u0438\u043c\u043e\u0441\u0442\u044c:<\/p>\n<pre><code class=\"java\">testImplementation \"org.apache.flink:flink-statebackend-rocksdb:${flinkVersion}\"<\/code><\/pre>\n<p>\u0422\u0435\u043f\u0435\u0440\u044c \u043e\u0441\u0442\u0430\u043b\u043e\u0441\u044c \u0442\u043e\u043b\u044c\u043a\u043e \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c \u044d\u0442\u043e\u0442 \u043a\u043e\u043d\u0444\u0438\u0433 \u0432 E2E-\u0442\u0435\u0441\u0442\u0430\u0445, \u0434\u043e\u0431\u0430\u0432\u0438\u0432 \u0435\u0433\u043e \u0432 \u0430\u043d\u043d\u043e\u0442\u0430\u0446\u0438\u044e @E2ETest. \u0412\u043e\u0442 \u0447\u0442\u043e \u0432 \u0438\u0442\u043e\u0433\u0435 \u0432\u044b \u0441\u043c\u043e\u0436\u0435\u0442\u0435 \u043d\u0430\u0431\u043b\u044e\u0434\u0430\u0442\u044c \u0432 \u043b\u043e\u0433\u0430\u0445:<\/p>\n<pre><code>INFO 35116 --- [ger-io-thread-1] o.a.flink.runtime.jobmaster.JobMaster    : Using job\/cluster config to configure application-defined state backend: EmbeddedRocksDBStateBackend{, localRocksDbDirectories=null, enableIncrementalCheckpointing=FALSE, numberOfTransferThreads=-1, writeBatchSize=-1}<\/code><\/pre>\n<p><strong>\u0412\u0430\u0436\u043d\u043e\u0435 \u0437\u0430\u043c\u0435\u0447\u0430\u043d\u0438\u0435. <\/strong>\u0412 \u043f\u0440\u043e\u0446\u0435\u0441\u0441\u0435 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u043d\u0438\u044f RocksDb \u043e\u043d \u0441\u043e\u0437\u0434\u0430\u0451\u0442 \u0441\u0432\u043e\u0438 \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u044b\u0435 \u0434\u0438\u0440\u0435\u043a\u0442\u043e\u0440\u0438\u0438 \u0441 \u0434\u043e\u0441\u0442\u0430\u0442\u043e\u0447\u043d\u043e \u0434\u043b\u0438\u043d\u043d\u044b\u043c\u0438 \u043f\u0443\u0442\u044f\u043c\u0438. \u041d\u0430\u043f\u0440\u0438\u043c\u0435\u0440, \u0434\u043b\u044f Mac:<\/p>\n<pre><code>INFO 35116 --- [essages (1\/2)#0] .f.c.s.s.RocksDBKeyedStateBackendBuilder : Finished building RocksDB keyed state-backend at \/var\/folders\/_y\/gd8sxnq91z9glrxjlkrj98tnbsxn37\/T\/junit7312928894650170068\/junit1795779582692257025\/minicluster_aabeebf82ae121d1fa503365b6aa7eb7\/tm_0\/tmp\/job_9a5c9344174adde1453cf361f9a0f43f_op_StreamFlatMap_371a51a50a977e59af86fcb074c97b9f__1_2__uuid_270162db-c22d-49ff-bb62-9ed6d9ef2ac5.<\/code><\/pre>\n<p>\u041f\u043e\u044d\u0442\u043e\u043c\u0443 \u043f\u0440\u0438 \u0437\u0430\u043f\u0443\u0441\u043a\u0435 \u043a\u043e\u0434\u0430 \u043d\u0430 Windows \u0432\u044b \u043c\u043e\u0436\u0435\u0442\u0435 \u043f\u043e\u043b\u0443\u0447\u0438\u0442\u044c \u043d\u0435\u043e\u0436\u0438\u0434\u0430\u043d\u043d\u0443\u044e \u043e\u0448\u0438\u0431\u043a\u0443:<\/p>\n<pre><code class=\"java\">Caused by: java.io.IOException: The directory path length (275) is longer than the directory path length limit for Windows (247): C:\\Users\\User\\AppData\\Local\\Temp\\junit15245919089769733580\\junit7040907885113472950\\minicluster_7d9295c691c345c6209f2bb0db16b593\\tm_0\\tmp\\job_7e441850886077eb74921aa5bb0e41f4_op_StreamFlatMap_371a51a50a977e59af86fcb074c97b9f__1_2__uuid_bcc33f69-17b2-4424-9620-25922af3be34\\db at org.apache.flink.contrib.streaming.state.RocksDBOperationUtils.throwExceptionIfPathLengthExceededOnWindows(RocksDBOperationUtils.java:285) ~[flink-statebackend-rocksdb-1.17.0.jar:1.17.0] at org.apache.flink.contrib.streaming.state.RocksDBOperationUtils.openDB(RocksDBOperationUtils.java:85) ~[flink-statebackend-rocksdb-1.17.0.jar:1.17.0] at org.apache.flink.contrib.streaming.state.restore.RocksDBHandle.loadDb(RocksDBHandle.java:134) ~[flink-statebackend-rocksdb-1.17.0.jar:1.17.0] at org.apache.flink.contrib.streaming.state.restore.RocksDBHandle.openDB(RocksDBHandle.java:113) ~[flink-statebackend-rocksdb-1.17.0.jar:1.17.0] at org.apache.flink.contrib.streaming.state.restore.RocksDBNoneRestoreOperation.restore(RocksDBNoneRestoreOperation.java:62) ~[flink-statebackend-rocksdb-1.17.0.jar:1.17.0] at org.apache.flink.contrib.streaming.state.RocksDBKeyedStateBackendBuilder.build(RocksDBKeyedStateBackendBuilder.java:325) ~[flink-statebackend-rocksdb-1.17.0.jar:1.17.0] ... 18 common frames omitted Caused by: org.rocksdb.RocksDBException: Failed to create a directory: C:\\Users\\User\\AppData\\Local\\Temp\\junit15245919089769733580\\junit7040907885113472950\\minicluster_7d9295c691c345c6209f2bb0db16b593\\tm_0\\tmp\\job_7e441850886077eb74921aa5bb0e41f4_op_StreamFlatMap_371a51a50a977e59af86fcb074c97b9f__1_2__uuid_bcc33f69-17b2-4424-9620-25922af3be34\\db: The directory path length (275) is longer than the directory path length limit for Windows (247). at org.rocksdb.RocksDB.open(Native Method) ~[frocksdbjni-6.20.3-ververica-2.0.jar:na] at org.rocksdb.RocksDB.open(RocksDB.java:306) ~[frocksdbjni-6.20.3-ververica-2.0.jar:na] at org.apache.flink.contrib.streaming.state.RocksDBOperationUtils.openDB(RocksDBOperationUtils.java:75) ~[flink-statebackend-rocksdb-1.17.0.jar:1.17.0] ... 22 common frames omitted<\/code><\/pre>\n<p>\u0415\u0441\u0442\u044c \u0441\u0442\u0430\u0442\u044c\u0438 \u043e \u0442\u043e\u043c, \u043a\u0430\u043a \u0440\u0435\u0448\u0438\u0442\u044c \u044d\u0442\u0443 \u043f\u0440\u043e\u0431\u043b\u0435\u043c\u0443. \u041d\u043e \u0443 \u043d\u0435\u043a\u043e\u0442\u043e\u0440\u044b\u0445 \u0440\u0430\u0437\u0440\u0430\u0431\u043e\u0442\u0447\u0438\u043a\u043e\u0432 \u043e\u043f\u0438\u0441\u0430\u043d\u043d\u044b\u0435 \u0440\u0435\u0448\u0435\u043d\u0438\u044f \u043d\u0435 \u0441\u0440\u0430\u0431\u043e\u0442\u0430\u043b\u0438. \u041f\u043e\u044d\u0442\u043e\u043c\u0443 \u043f\u0440\u0438 \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0438 Flink-\u0437\u0430\u0434\u0430\u0447 \u0441 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u043d\u0438\u0435\u043c embedded RocksDB \u0440\u0435\u043a\u043e\u043c\u0435\u043d\u0434\u0443\u0435\u0442\u0441\u044f \u043d\u0435 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c Windows.<\/p>\n<p><a class=\"anchor\" name=\"10\" id=\"10\"><\/a><\/p>\n<h2>\u0412\u044b\u0432\u043e\u0434<\/h2>\n<p>\u041c\u044b \u0440\u0430\u0437\u043e\u0431\u0440\u0430\u043b\u0438, \u043a\u0430\u043a \u043d\u0430\u043f\u0438\u0441\u0430\u0442\u044c E2E-\u0442\u0435\u0441\u0442 \u043d\u0430 Flink \u0438 Spring-\u0434\u0436\u043e\u0431\u0443 \u0441 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u043d\u0438\u0435\u043c Kafka \u0438 Testcontainers. \u041c\u044b \u0441\u043e\u0437\u0434\u0430\u043b\u0438 \u0443\u0434\u043e\u0431\u043d\u044b\u0435 \u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u0438 \u0438 \u043f\u043e\u0433\u043e\u0432\u043e\u0440\u0438\u043b\u0438 \u043e\u0431 \u043e\u0441\u043d\u043e\u0432\u043d\u044b\u0445 \u043f\u0440\u043e\u0431\u043b\u0435\u043c\u0430\u0445, \u043a\u043e\u0442\u043e\u0440\u044b\u0435 \u043d\u0443\u0436\u043d\u043e \u0443\u0447\u0438\u0442\u044b\u0432\u0430\u0442\u044c \u043f\u0440\u0438 \u0442\u0430\u043a\u043e\u043c \u0432\u0438\u0434\u0435 \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f.<\/p>\n<p>\u041d\u0430 \u044d\u0442\u043e\u043c \u0441\u0435\u0440\u0438\u044e \u0441\u0442\u0430\u0442\u0435\u0439 \u043f\u0440\u043e \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435 \u044f \u0437\u0430\u043a\u043e\u043d\u0447\u0443. \u0414\u0430\u043b\u044c\u0448\u0435 \u043f\u0435\u0440\u0435\u0439\u0434\u0443 \u043a \u0440\u0430\u0437\u0431\u043e\u0440\u0443 \u0442\u0430\u0439\u043c\u0435\u0440\u043e\u0432 \u0432 Flink \u0441 \u0431\u043e\u043b\u0435\u0435 \u0441\u043b\u043e\u0436\u043d\u044b\u043c \u0445\u0440\u0430\u043d\u0435\u043d\u0438\u0435\u043c \u0441\u043e\u0441\u0442\u043e\u044f\u043d\u0438\u044f. \u041c\u044b \u043f\u043e\u0441\u043c\u043e\u0442\u0440\u0438\u043c, \u043a\u0430\u043a \u043e\u0442\u043b\u043e\u0436\u0438\u0442\u044c \u043a\u0430\u043a\u043e\u0435-\u0442\u043e \u0434\u0435\u0439\u0441\u0442\u0432\u0438\u0435 \u0438 \u043e\u0442\u043f\u0440\u0430\u0432\u0438\u0442\u044c \u0441\u043e\u0431\u044b\u0442\u0438\u0435 \u043f\u043e \u0443\u0441\u0442\u0430\u043d\u043e\u0432\u043b\u0435\u043d\u043d\u043e\u043c\u0443 \u0442\u0430\u0439\u043c\u0435\u0440\u0443 \u0441 \u043f\u043e\u043c\u043e\u0449\u044c\u044e \u0441\u043e\u0441\u0442\u043e\u044f\u043d\u0438\u044f. \u0420\u0430\u0437\u0431\u0435\u0440\u0451\u043c, \u043a\u0430\u043a\u0438\u0435 \u043f\u0440\u043e\u0431\u043b\u0435\u043c\u044b \u043c\u043e\u0433\u0443\u0442 \u0432\u043e\u0437\u043d\u0438\u043a\u043d\u0443\u0442\u044c \u043f\u0440\u0438 \u0431\u043e\u043b\u0435\u0435 \u0441\u043b\u043e\u0436\u043d\u043e\u043c \u0432\u0437\u0430\u0438\u043c\u043e\u0434\u0435\u0439\u0441\u0442\u0432\u0438\u0438 \u0441 \u0441\u043e\u0441\u0442\u043e\u044f\u043d\u0438\u044f\u043c\u0438. \u041a\u043e\u043d\u0435\u0447\u043d\u043e, \u043f\u0440\u043e\u0434\u043e\u043b\u0436\u0438\u043c \u043f\u043e\u043a\u0440\u044b\u0432\u0430\u0442\u044c \u043a\u043e\u0434 \u0442\u0435\u0441\u0442\u0430\u043c\u0438 \u2014 \u0430 \u0437\u043d\u0430\u0447\u0438\u0442, \u0432\u0441\u0442\u0440\u0435\u0442\u0438\u043c\u0441\u044f \u0441 \u043d\u043e\u0432\u044b\u043c\u0438 \u043c\u0435\u0445\u0430\u043d\u0438\u0437\u043c\u0430\u043c\u0438 \u0438 \u043f\u0440\u0430\u043a\u0442\u0438\u043a\u0430\u043c\u0438 \u043f\u0440\u0438 \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0438.<\/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\/819681\/\"> https:\/\/habr.com\/ru\/articles\/819681\/<\/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>\u041f\u0440\u0438\u0432\u0435\u0442, \u0425\u0430\u0431\u0440! \u0421 \u0432\u0430\u043c\u0438 \u0410\u043b\u0435\u043a\u0441\u0430\u043d\u0434\u0440 \u0411\u043e\u0431\u0440\u044f\u043a\u043e\u0432, \u0442\u0435\u0445\u043b\u0438\u0434 \u0432 \u043a\u043e\u043c\u0430\u043d\u0434\u0435 \u041c\u0422\u0421 \u0410\u043d\u0430\u043b\u0438\u0442\u0438\u043a\u0438. \u042f \u043a \u0432\u0430\u043c \u0441 \u043d\u043e\u0432\u043e\u0439 \u0441\u0442\u0430\u0442\u044c\u0451\u0439 \u0438\u0437 \u0446\u0438\u043a\u043b\u0430 \u043f\u0440\u043e \u0444\u0440\u0435\u0439\u043c\u0432\u043e\u0440\u043a Apache Flink. <\/p>\n<p>\u0412 <a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/812905\/\">\u043f\u0440\u0435\u0434\u044b\u0434\u0443\u0449\u0435\u0439 \u0447\u0430\u0441\u0442\u0438<\/a> \u044f \u0440\u0430\u0441\u0441\u043a\u0430\u0437\u0430\u043b, \u043a\u0430\u043a \u0441\u043e\u0437\u0434\u0430\u0442\u044c Unit-\u0442\u0435\u0441\u0442 \u043d\u0430 \u043f\u043e\u043b\u043d\u043e\u0446\u0435\u043d\u043d\u0443\u044e \u0434\u0436\u043e\u0431\u0443 Flink \u0438 \u043e\u0442\u0434\u0435\u043b\u044c\u043d\u044b\u0435 stateful-\u043e\u043f\u0435\u0440\u0430\u0442\u043e\u0440\u044b \u0441 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u043d\u0438\u0435\u043c Flink MiniCluster. \u0415\u0449\u0451 \u043c\u044b \u043d\u0430\u0443\u0447\u0438\u043b\u0438\u0441\u044c \u0437\u0430\u043f\u0443\u0441\u043a\u0430\u0442\u044c \u043c\u0438\u043d\u0438-\u043a\u043b\u0430\u0441\u0442\u0435\u0440 \u043e\u0434\u0438\u043d \u0440\u0430\u0437 \u043f\u0435\u0440\u0435\u0434 \u0432\u0441\u0435\u043c\u0438 \u0442\u0435\u0441\u0442\u043e\u0432\u044b\u043c\u0438 \u043a\u043b\u0430\u0441\u0441\u0430\u043c\u0438, \u043a\u043e\u0442\u043e\u0440\u044b\u0435 \u043d\u0443\u0436\u0434\u0430\u044e\u0442\u0441\u044f \u0432 \u043d\u0451\u043c. \u0412 \u0434\u043e\u043f\u043e\u043b\u043d\u0435\u043d\u0438\u0435 \u0441\u043e\u0437\u0434\u0430\u043b\u0438 \u0432\u0441\u043f\u043e\u043c\u043e\u0433\u0430\u0442\u0435\u043b\u044c\u043d\u044b\u0435 \u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u0438 \u0438 \u0430\u043d\u043d\u043e\u0442\u0430\u0446\u0438\u0438, \u0437\u043d\u0430\u0447\u0438\u0442\u0435\u043b\u044c\u043d\u043e \u0440\u0430\u0437\u0434\u0435\u043b\u044f\u044f \u043e\u0442\u0432\u0435\u0442\u0441\u0442\u0432\u0435\u043d\u043d\u043e\u0441\u0442\u044c \u0432 \u0442\u0435\u0441\u0442\u0430\u0445 \u0438 \u0443\u043f\u0440\u043e\u0449\u0430\u044f \u043b\u043e\u0433\u0438\u043a\u0443 \u043d\u0430\u043f\u0438\u0441\u0430\u043d\u0438\u044f \u043d\u043e\u0432\u044b\u0445 \u0442\u0435\u0441\u0442\u043e\u0432. <\/p>\n<p>\u0412 \u043f\u0440\u0435\u0434\u044b\u0434\u0443\u0449\u0438\u0445 \u0442\u0435\u0441\u0442\u0430\u0445 \u043d\u0430 \u0434\u0436\u043e\u0431\u0443 \u043c\u044b \u043d\u0435 \u0437\u0430\u0442\u0440\u0430\u0433\u0438\u0432\u0430\u043b\u0438 \u0438\u043d\u0442\u0435\u0433\u0440\u0430\u0446\u0438\u044e \u0441 Kafka, \u0432\u0435\u0434\u044c \u043d\u0430\u043c \u0431\u044b\u043b\u0438 \u043d\u0435 \u0432\u0430\u0436\u043d\u044b \u0440\u0435\u0430\u043b\u044c\u043d\u044b\u0435 source \u0438 sink. \u0412 \u044d\u0442\u043e\u0439 \u0441\u0442\u0430\u0442\u044c\u0435 \u043f\u0440\u043e\u0434\u043e\u043b\u0436\u0438\u043c \u0440\u0430\u0437\u0431\u0438\u0440\u0430\u0442\u044c\u0441\u044f \u0432 \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0438 \u0438 \u043d\u0430\u043f\u0438\u0448\u0435\u043c \u043f\u043e\u043b\u043d\u043e\u0446\u0435\u043d\u043d\u044b\u0439 E2E-\u0442\u0435\u0441\u0442, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043e\u0445\u0432\u0430\u0442\u0438\u0442 Kafka \u0438 Flink \u0432\u043c\u0435\u0441\u0442\u0435 \u0441 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u043d\u0438\u0435\u043c Testcontainers. \u0422\u0430\u043a\u0436\u0435 \u0440\u0430\u0441\u0441\u043c\u043e\u0442\u0440\u0438\u043c \u043d\u0435\u043e\u0447\u0435\u0432\u0438\u0434\u043d\u044b\u0435 \u043f\u0440\u043e\u0431\u043b\u0435\u043c\u044b \u0432 \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0438 \u0438 \u043d\u043e\u0432\u044b\u0435 \u0443\u043d\u0438\u0432\u0435\u0440\u0441\u0430\u043b\u044c\u043d\u044b\u0435 \u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u0438.<\/p>\n<figure class=\"full-width\"><\/figure>\n<details class=\"spoiler\">\n<summary>\u0421\u043f\u0438\u0441\u043e\u043a \u043c\u043e\u0438\u0445 \u043f\u043e\u0441\u0442\u043e\u0432 \u043f\u0440\u043e Flink<\/summary>\n<div class=\"spoiler__content\">\n<ol>\n<li>\n<p><a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/772898\/\">\u0412\u0432\u0435\u0434\u0435\u043d\u0438\u0435 \u0432 Apache Flink: \u043e\u0441\u0432\u0430\u0438\u0432\u0430\u0435\u043c \u0444\u0440\u0435\u0439\u043c\u0432\u043e\u0440\u043a \u043d\u0430 \u0440\u0435\u0430\u043b\u044c\u043d\u044b\u0445 \u043f\u0440\u0438\u043c\u0435\u0440\u0430\u0445<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/775970\/\">\u041f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0435 \u043f\u043e\u0434 Apache Flink: \u0441 \u0447\u0435\u0433\u043e \u043d\u0430\u0447\u0430\u0442\u044c?<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/786012\/\">Apache Flink. \u041a\u0430\u043a \u0440\u0430\u0431\u043e\u0442\u0430\u0435\u0442 \u0434\u0435\u0434\u0443\u043f\u043b\u0438\u043a\u0430\u0446\u0438\u044f \u0434\u0430\u043d\u043d\u044b\u0445 \u0432 \u043f\u043e\u0442\u043e\u043a\u0435 Kafka-to-Kafka?<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/798667\/\">\u0414\u0438\u043d\u0430\u043c\u0438\u0447\u0435\u0441\u043a\u043e\u0435 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u0438\u0435 \u0432\u044b\u0445\u043e\u0434\u043d\u043e\u0433\u043e \u0442\u043e\u043f\u0438\u043a\u0430<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/801693\/\">\u041a\u0430\u043a \u043f\u0440\u043e\u0432\u0435\u0441\u0442\u0438 unit-\u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435 Flink-\u043e\u043f\u0435\u0440\u0430\u0442\u043e\u0440\u043e\u0432: Test Harness<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/812905\/\">Unit-\u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435 Flink-\u043e\u043f\u0435\u0440\u0430\u0442\u043e\u0440\u043e\u0432, Job: Flink MiniCluster<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"https:\/\/habr.com\/ru\/companies\/ru_mts\/articles\/819681\/\"><strong>E2E-\u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435 Flink Job \u0441 Kafka<\/strong><\/a><\/p>\n<\/li>\n<\/ol>\n<\/div>\n<\/details>\n<p>\u0412\u0435\u0441\u044c \u0440\u0430\u0437\u0431\u0438\u0440\u0430\u0435\u043c\u044b\u0439 \u0438\u0441\u0445\u043e\u0434\u043d\u044b\u0439 \u043a\u043e\u0434 \u043c\u043e\u0436\u043d\u043e \u043d\u0430\u0439\u0442\u0438 \u0432 \u0440\u0435\u043f\u043e\u0437\u0438\u0442\u043e\u0440\u0438\u0438 <a href=\"https:\/\/github.com\/AlexanderBobryakov\/flink-spring\">AlexanderBobryakov\/flink-spring<\/a>. \u0412 master-\u0432\u0435\u0442\u043a\u0435 \u2014 \u0438\u0442\u043e\u0433\u043e\u0432\u044b\u0439 \u043f\u0440\u043e\u0435\u043a\u0442 \u043f\u043e \u0432\u0441\u0435\u0439 \u0441\u0435\u0440\u0438\u0438 \u0441\u0442\u0430\u0442\u0435\u0439. \u042d\u0442\u0430 \u0441\u0442\u0430\u0442\u044c\u044f \u0441\u043e\u043e\u0442\u0432\u0435\u0442\u0441\u0442\u0432\u0443\u0435\u0442 \u0440\u0435\u043b\u0438\u0437\u043d\u043e\u0439 \u0432\u0435\u0442\u043a\u0435 \u043f\u043e\u0434 \u043d\u0430\u0437\u0432\u0430\u043d\u0438\u0435\u043c <a href=\"https:\/\/github.com\/AlexanderBobryakov\/flink-spring\/tree\/release\/6_e2e_deduplicator_test\">release\/6_e2e_deduplicator_test<\/a>.<\/p>\n<details class=\"spoiler\">\n<summary>\u041e\u0433\u043b\u0430\u0432\u043b\u0435\u043d\u0438\u0435 \u0441\u0442\u0430\u0442\u044c\u0438<\/summary>\n<div class=\"spoiler__content\">\n<ol>\n<li>\n<p><a href=\"#1\">E2E-\u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#2\">\u041f\u043e\u0434\u043d\u0438\u043c\u0430\u0435\u043c Kafka \u0441 \u043f\u043e\u043c\u043e\u0449\u044c\u044e Testcontainers<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#3\">\u0410\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u0438 \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u0434\u043b\u044f Kafka<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#4\">KafkaTestConsumer<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#5\">TestKafkaFacade<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#6\">KafkaTopicCreatorConfig<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#7\">E2E-\u0442\u0435\u0441\u0442 \u043d\u0430 Flink Job<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#8\">\u0411\u0435\u0437\u043e\u043f\u0430\u0441\u043d\u043e\u0435 \u0437\u0430\u0432\u0435\u0440\u0448\u0435\u043d\u0438\u0435 E2E-\u0442\u0435\u0441\u0442\u043e\u0432<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#9\">RocksDB \u0432 E2E-\u0442\u0435\u0441\u0442\u0430\u0445<\/a><\/p>\n<\/li>\n<li>\n<p><a href=\"#10\">\u0412\u044b\u0432\u043e\u0434<\/a><\/p>\n<\/li>\n<\/ol>\n<\/div>\n<\/details>\n<p><a class=\"anchor\" name=\"1\" id=\"1\"><\/a><\/p>\n<h2>E2E-\u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435<\/h2>\n<p>E2E-\u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435 \u043e\u0445\u0432\u0430\u0442\u044b\u0432\u0430\u0435\u0442 \u043f\u043e\u0432\u0435\u0434\u0435\u043d\u0438\u0435 \u0432\u0441\u0435\u0439 \u0441\u0438\u0441\u0442\u0435\u043c\u044b \u043e\u0442 \u043d\u0430\u0447\u0430\u043b\u0430 \u0434\u043e \u043a\u043e\u043d\u0446\u0430. \u041f\u043e \u0441\u043f\u0435\u0446\u0438\u0444\u0438\u043a\u0435 \u043d\u0430\u0448\u0435\u0439 \u0434\u0436\u043e\u0431\u044b, \u043a\u043e\u0442\u043e\u0440\u0443\u044e \u043c\u044b \u0440\u0430\u0441\u0441\u043c\u0430\u0442\u0440\u0438\u0432\u0430\u043b\u0438 \u0432 \u043f\u0440\u0435\u0434\u044b\u0434\u0443\u0449\u0438\u0445 \u0447\u0430\u0441\u0442\u044f\u0445, \u0432 \u043d\u0430\u0447\u0430\u043b\u0435 \u0435\u0441\u0442\u044c Kafka-\u0442\u043e\u043f\u0438\u043a \u0441 \u0434\u0430\u043d\u043d\u044b\u043c\u0438 ClickMessage, \u0430 \u043d\u0430 \u0432\u044b\u0445\u043e\u0434\u0435 \u2014 \u043c\u043d\u043e\u0433\u043e \u0440\u0430\u0437\u043d\u044b\u0445 product-\u0442\u043e\u043f\u0438\u043a\u043e\u0432. \u0417\u043d\u0430\u0447\u0438\u0442, \u0432 \u0442\u0435\u0441\u0442\u0435 \u0432\u0441\u0451 \u044d\u0442\u043e \u0434\u043e\u043b\u0436\u043d\u043e \u0443\u0447\u0438\u0442\u044b\u0432\u0430\u0442\u044c\u0441\u044f. <\/p>\n<p>\u0426\u0435\u043b\u044c \u0442\u0430\u043a\u043e\u0433\u043e \u0442\u0435\u0441\u0442\u0430 \u2014 \u043f\u0440\u043e\u0432\u0435\u0440\u0438\u0442\u044c \u0432\u0441\u0435 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u043c\u044b\u0435 \u0438\u043d\u0442\u0435\u0433\u0440\u0430\u0446\u0438\u0438 \u043c\u0435\u0436\u0434\u0443 \u0441\u043e\u0431\u043e\u0439. \u0418\u043d\u0430\u0447\u0435 \u0431\u043b\u043e\u043a\u0438 \u0441\u0438\u0441\u0442\u0435\u043c\u044b \u043c\u043e\u0433\u0443\u0442 \u0440\u0430\u0431\u043e\u0442\u0430\u0442\u044c \u043f\u043e \u043e\u0442\u0434\u0435\u043b\u044c\u043d\u043e\u0441\u0442\u0438, \u0430 \u0432\u043c\u0435\u0441\u0442\u0435 \u0432\u0441\u0451 \u0441\u043b\u043e\u043c\u0430\u0435\u0442\u0441\u044f. \u0412 \u043d\u0430\u0448\u0435\u043c \u0441\u043b\u0443\u0447\u0430\u0435 \u0434\u043e\u043b\u0436\u0435\u043d \u043f\u043e\u0434\u043d\u0438\u043c\u0430\u0442\u044c\u0441\u044f \u0432\u0435\u0441\u044c Spring-\u043a\u043e\u043d\u0442\u0435\u043a\u0441\u0442, \u0441\u0442\u0430\u0440\u0442\u043e\u0432\u0430\u0442\u044c Flink \u0438 Kafka, \u043a\u0430\u043a \u0431\u0443\u0434\u0442\u043e \u043c\u044b \u0437\u0430\u043f\u0443\u0441\u043a\u0430\u0435\u043c \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0435 \u043d\u0430 \u043f\u0440\u043e\u0434\u0435.<\/p>\n<p><a class=\"anchor\" name=\"2\" id=\"2\"><\/a><\/p>\n<h2>\u041f\u043e\u0434\u043d\u0438\u043c\u0430\u0435\u043c Kafka \u0441 \u043f\u043e\u043c\u043e\u0449\u044c\u044e Testcontainers<\/h2>\n<p><a href=\"https:\/\/testcontainers.com\">Testcontainers<\/a> \u2014 \u200b\u200b\u044d\u0442\u043e \u0431\u0438\u0431\u043b\u0438\u043e\u0442\u0435\u043a\u0430 Java, \u043a\u043e\u0442\u043e\u0440\u0430\u044f \u043f\u043e\u0434\u0434\u0435\u0440\u0436\u0438\u0432\u0430\u0435\u0442 \u0442\u0435\u0441\u0442\u044b JUnit. \u041e\u043d\u0430 \u0434\u0430\u0451\u0442 \u0432\u043e\u0437\u043c\u043e\u0436\u043d\u043e\u0441\u0442\u044c \u0437\u0430\u043f\u0443\u0441\u043a\u0430\u0442\u044c \u0432 \u043d\u0438\u0445 \u0432\u0441\u0451, \u0447\u0442\u043e \u043c\u043e\u0436\u0435\u0442 \u0437\u0430\u043f\u0443\u0441\u043a\u0430\u0442\u044c\u0441\u044f \u0432 Docker. \u0417\u043d\u0430\u0447\u0438\u0442, \u0432\u044b \u043c\u043e\u0436\u0435\u0442\u0435 \u043f\u0440\u043e\u0432\u0435\u0440\u0438\u0442\u044c \u043b\u044e\u0431\u0443\u044e \u0438\u043d\u0442\u0435\u0433\u0440\u0430\u0446\u0438\u044e \u0432\u0430\u0448\u0435\u0433\u043e \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u044f: \u0441 \u0411\u0414, \u0431\u0440\u043e\u043a\u0435\u0440\u0430\u043c\u0438 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0439, \u0434\u0440\u0443\u0433\u0438\u043c\u0438 \u0441\u0435\u0440\u0432\u0438\u0441\u0430\u043c\u0438 \u0438 \u0442\u0430\u043a \u0434\u0430\u043b\u0435\u0435. <\/p>\n<p>\u0421\u0446\u0435\u043d\u0430\u0440\u0438\u0439 \u043d\u0430\u043f\u0438\u0441\u0430\u043d\u0438\u044f \u0442\u0435\u0441\u0442\u0430 \u0432 \u0438\u0442\u043e\u0433\u0435 \u0432\u044b\u0433\u043b\u044f\u0434\u0438\u0442 \u0442\u0430\u043a:<\/p>\n<ol>\n<li>\n<p>\u041e\u043f\u0440\u0435\u0434\u0435\u043b\u0438\u0442\u044c \u0432 \u0442\u0435\u0441\u0442\u0435 Testcontainers \u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440 \u2014 \u043d\u0430\u043f\u0440\u0438\u043c\u0435\u0440, \u0434\u043b\u044f Kafka.<\/p>\n<\/li>\n<li>\n<p>\u0417\u0430\u043f\u0443\u0441\u0442\u0438\u0442\u044c Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440.<\/p>\n<\/li>\n<li>\n<p>\u041f\u0440\u043e\u0431\u0440\u043e\u0441\u0438\u0442\u044c \u0441\u0432\u043e\u0439\u0441\u0442\u0432\u0430 \u0434\u043b\u044f \u043f\u043e\u0434\u043a\u043b\u044e\u0447\u0435\u043d\u0438\u044f \u043a Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440\u0443 \u0432 \u043a\u043e\u043d\u0444\u0438\u0433 \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u044f.<\/p>\n<\/li>\n<li>\n<p>\u0417\u0430\u043f\u0443\u0441\u0442\u0438\u0442\u044c \u0442\u0435\u0441\u0442, \u0432 \u043a\u043e\u0442\u043e\u0440\u043e\u043c \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0435 \u043f\u043e\u0434\u043a\u043b\u044e\u0447\u0430\u0435\u0442\u0441\u044f \u0441\u043e\u0433\u043b\u0430\u0441\u043d\u043e \u043a\u043e\u043d\u0444\u0438\u0433\u0443 \u043a Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440\u0443.<\/p>\n<\/li>\n<\/ol>\n<p>\u0412 \u043f\u0440\u0435\u0434\u044b\u0434\u0443\u0449\u0435\u0439 \u0441\u0442\u0430\u0442\u044c\u0435 \u043c\u044b \u0437\u0430\u0442\u0440\u043e\u043d\u0443\u043b\u0438 \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u0435 \u043a\u0430\u0441\u0442\u043e\u043c\u043d\u044b\u0445 \u0430\u043d\u043d\u043e\u0442\u0430\u0446\u0438\u0439 \u043f\u043e\u0434 \u0442\u0435\u0441\u0442\u044b. \u042d\u0442\u043e \u0431\u044b\u043b\u043e \u0434\u043e\u0441\u0442\u0430\u0442\u043e\u0447\u043d\u043e \u0443\u0434\u043e\u0431\u043d\u043e, \u043f\u043e\u044d\u0442\u043e\u043c\u0443 \u043f\u0440\u0435\u0434\u043b\u0430\u0433\u0430\u044e \u043f\u0440\u0438\u0434\u0435\u0440\u0436\u0438\u0432\u0430\u0442\u044c\u0441\u044f \u0430\u043d\u0430\u043b\u043e\u0433\u0438\u0447\u043d\u043e\u0433\u043e \u043f\u043e\u0434\u0445\u043e\u0434\u0430 \u0438 \u0432 \u044d\u0442\u043e\u0442 \u0440\u0430\u0437. \u0414\u043b\u044f \u043d\u0430\u0447\u0430\u043b\u0430 \u043d\u0430\u043c \u043d\u0443\u0436\u043d\u043e \u043f\u043e\u0434\u043a\u043b\u044e\u0447\u0435\u043d\u0438\u0435 \u043a Kafka. TestContainers \u043f\u0440\u0435\u0434\u043e\u0441\u0442\u0430\u0432\u043b\u044f\u0435\u0442 Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440 <a href=\"https:\/\/www.testcontainers.org\/modules\/kafka\/\">\u00ab\u0438\u0437 \u043a\u043e\u0440\u043e\u0431\u043a\u0438\u00bb<\/a>, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043d\u0435 \u043d\u0443\u0436\u0434\u0430\u0435\u0442\u0441\u044f \u0432 <strong>\u043e\u0442\u0434\u0435\u043b\u044c\u043d\u043e\u0439<\/strong> \u0438\u043d\u0438\u0446\u0438\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438 Zookeeper. \u0415\u0449\u0451 TestContainers \u043c\u043e\u0436\u043d\u043e \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c, \u0447\u0442\u043e\u0431\u044b \u0441\u043e\u0437\u0434\u0430\u0442\u044c \u0432\u0441\u0435 \u043d\u0435\u043e\u0431\u0445\u043e\u0434\u0438\u043c\u044b\u0435 \u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440\u044b \u043d\u0430 \u043e\u0441\u043d\u043e\u0432\u0435 \u043b\u044e\u0431\u044b\u0445 \u0434\u043e\u043a\u0435\u0440-\u043e\u0431\u0440\u0430\u0437\u043e\u0432.<\/p>\n<p>\u0414\u043b\u044f \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u043d\u0438\u044f \u0433\u043e\u0442\u043e\u0432\u043e\u0433\u043e Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440\u0430 \u043c\u043e\u0436\u0435\u043c \u0432\u043e\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c\u0441\u044f \u0437\u0430\u0432\u0438\u0441\u0438\u043c\u043e\u0441\u0442\u044c\u044e:<\/p>\n<pre><code>testImplementation \"org.testcontainers:kafka\"<\/code><\/pre>\n<p>\u0410\u043d\u0430\u043b\u043e\u0433\u0438\u0447\u043d\u043e \u0442\u043e\u043c\u0443, \u043a\u0430\u043a \u043c\u044b \u0441\u043e\u0437\u0434\u0430\u043b\u0438 JUnit Extension \u0434\u043b\u044f \u0441\u0442\u0430\u0440\u0442\u0430 Flink MiniCluster, \u0441\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u043d\u043e\u0432\u044b\u0439 Extension \u0434\u043b\u044f \u0441\u0442\u0430\u0440\u0442\u0430 Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440\u0430:<\/p>\n<pre><code class=\"java\">@Slf4j @SuppressWarnings({\"PMD.AvoidUsingVolatile\"}) public class KafkaContainerExtension implements BeforeAllCallback, ExtensionContext.Store.CloseableResource {    private static final KafkaContainer KAFKA =        new KafkaContainer(DockerImageName.parse(\"confluentinc\/cp-kafka:7.3.2\"))            .withEnv(\"KAFKA_AUTO_CREATE_TOPICS_ENABLE\", \"false\");     private static final Lock LOCK = new ReentrantLock();    private static volatile boolean started;     @Override    public void beforeAll(ExtensionContext context) {        LOCK.lock();        try {            if (!started) {                log.info(\"Start Kafka Container\");                started = true;                Startables.deepStart(KAFKA).join();                System.setProperty(\"spring.kafka.bootstrap-servers\", KAFKA.getBootstrapServers());                System.setProperty(\"kafka.bootstrap-servers\", KAFKA.getBootstrapServers());                System.setProperty(\"spring.kafka.consumer.group-id\", \"group-id-spring\");                context.getRoot().getStore(GLOBAL).put(\"Kafka Container\", this);            }        } finally {            LOCK.unlock();        }    }     @Override    public void close() {        log.info(\"Close Kafka Container\");        KAFKA.close();        started = false;    } }<\/code><\/pre>\n<p>\u041a\u043e\u0434 \u043e\u0447\u0435\u043d\u044c \u043f\u043e\u0445\u043e\u0436 \u043d\u0430 \u043d\u0430\u0448 \u0441\u0443\u0449\u0435\u0441\u0442\u0432\u0443\u044e\u0449\u0438\u0439 FlinkClusterExtension, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u044f \u043e\u043f\u0438\u0441\u0430\u043b \u0432 \u043f\u0440\u043e\u0448\u043b\u043e\u0439 \u0441\u0442\u0430\u0442\u044c\u0435: \u043c\u044b \u0438\u043d\u0438\u0446\u0438\u0430\u043b\u0438\u0437\u0438\u0440\u0443\u0435\u043c Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440 \u043f\u043e \u0443\u043a\u0430\u0437\u0430\u043d\u043d\u043e\u043c\u0443 \u0434\u043e\u043a\u0435\u0440-\u043e\u0431\u0440\u0430\u0437\u0443, \u043f\u043e\u0442\u043e\u043c \u0432 beforeAll() \u0441\u0438\u043d\u0445\u0440\u043e\u043d\u043d\u043e \u0437\u0430\u043f\u0443\u0441\u043a\u0430\u0435\u043c \u0435\u0433\u043e \u0447\u0435\u0440\u0435\u0437 \u0432\u044b\u0437\u043e\u0432 Startables.deepStart(KAFKA).join(), \u043e\u0431\u0432\u044f\u0437\u044b\u0432\u0430\u044f \u0431\u043b\u043e\u043a\u0438\u0440\u043e\u0432\u043a\u0430\u043c\u0438. \u0412 \u043a\u043e\u043d\u0446\u0435 \u0432\u044b\u043f\u043e\u043b\u043d\u0435\u043d\u0438\u044f \u0432\u0441\u0435\u0445 \u0437\u0430\u0432\u0438\u0441\u0438\u043c\u044b\u0445 \u0442\u0435\u0441\u0442\u043e\u0432 \u0437\u0430\u043a\u0440\u044b\u0432\u0430\u0435\u043c \u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440 \u0432 \u043a\u043e\u043b\u0431\u044d\u043a-\u043c\u0435\u0442\u043e\u0434\u0435 \u0436\u0438\u0437\u043d\u0435\u043d\u043d\u043e\u0433\u043e \u0446\u0438\u043a\u043b\u0430 JUnit \u0442\u0435\u0441\u0442\u043e\u0432 close().<\/p>\n<p>\u0412\u043e\u0437\u043d\u0438\u043a\u0430\u0435\u0442 \u0432\u043e\u043f\u0440\u043e\u0441: \u043a\u0430\u043a \u043d\u0430\u0448\u0435 Spring-\u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0435 \u043f\u0440\u0438 \u0441\u0442\u0430\u0440\u0442\u0435 \u0431\u0443\u0434\u0435\u0442 \u043f\u043e\u0434\u043a\u043b\u044e\u0447\u0430\u0442\u044c\u0441\u044f \u043a Kafka? \u0412\u0435\u0434\u044c \u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440 \u0437\u0430\u043f\u0443\u0441\u043a\u0430\u0435\u0442\u0441\u044f \u043d\u0430 \u0441\u043b\u0443\u0447\u0430\u0439\u043d\u043e\u043c \u0441\u0432\u043e\u0431\u043e\u0434\u043d\u043e\u043c \u043f\u043e\u0440\u0442\u0443. \u0414\u043b\u044f \u044d\u0442\u043e\u0433\u043e \u043c\u044b \u043f\u0435\u0440\u0435\u0434\u0430\u0451\u043c \u043d\u0430\u0441\u0442\u0440\u043e\u0439\u043a\u0438 \u0432 \u043f\u0435\u0440\u0435\u043c\u0435\u043d\u043d\u044b\u0435 \u043e\u043a\u0440\u0443\u0436\u0435\u043d\u0438\u044f \u0447\u0435\u0440\u0435\u0437 System.setProperty() \u0432 \u0441\u0442\u0430\u0442\u0438\u0447\u0435\u0441\u043a\u043e\u043c \u043a\u043e\u043d\u0442\u0435\u043a\u0441\u0442\u0435 \u043d\u0435\u043f\u043e\u0441\u0440\u0435\u0434\u0441\u0442\u0432\u0435\u043d\u043d\u043e \u043f\u0435\u0440\u0435\u0434 \u0441\u0442\u0430\u0440\u0442\u043e\u043c \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u044f. \u041d\u0430\u0441\u0442\u0440\u043e\u0439\u043a\u0438 \u043f\u0435\u0440\u0435\u0434\u0430\u0451\u043c \u0441\u043e\u0433\u043b\u0430\u0441\u043d\u043e \u0441\u0442\u0440\u0443\u043a\u0442\u0443\u0440\u0435 application.yml, \u043f\u043e\u0442\u043e\u043c\u0443 \u0447\u0442\u043e \u043e\u043d\u0438 \u0431\u0443\u0434\u0443\u0442 \u00ab\u043f\u0435\u0440\u0435\u0437\u0430\u0442\u0438\u0440\u0430\u0442\u044c\u0441\u044f\u00bb \u0438\u0437 \u0443\u043a\u0430\u0437\u0430\u043d\u043d\u044b\u0445 \u043f\u0435\u0440\u0435\u043c\u0435\u043d\u043d\u044b\u0445 \u043e\u043a\u0440\u0443\u0436\u0435\u043d\u0438\u044f \u2014 <a href=\"https:\/\/docs.spring.io\/spring-boot\/docs\/current\/reference\/html\/features.html#features.external-config.typesafe-configuration-properties.relaxed-binding.environment-variables\">\u0434\u043e\u043a\u0443\u043c\u0435\u043d\u0442\u0430\u0446\u0438\u044f<\/a>:<\/p>\n<pre><code class=\"java\">kafka:  group-id: group_id  bootstrap-servers: localhost:29092<\/code><\/pre>\n<p>\u0414\u043b\u044f Spring \u043f\u0435\u0440\u0435\u0434\u0430\u0451\u043c \u0441\u0432\u043e\u0439\u0441\u0442\u0432\u0430, \u0447\u0442\u043e\u0431\u044b \u0430\u0432\u0442\u043e\u043c\u0430\u0442\u0438\u0447\u0435\u0441\u043a\u0438 \u0441\u043e\u0437\u0434\u0430\u043b\u0438\u0441\u044c \u043e\u0441\u043d\u043e\u0432\u043d\u044b\u0435 Spring-\u0431\u0438\u043d\u044b \u0434\u043b\u044f \u0438\u043d\u0442\u0435\u0433\u0440\u0430\u0446\u0438\u0438 \u0441 Kafka. \u041d\u0430\u043f\u0440\u0438\u043c\u0435\u0440, KafkaTemplate \u2014 \u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u044f \u043d\u0430\u0434 Kafka Producer, \u043a\u043e\u0442\u043e\u0440\u0430\u044f \u0443\u043c\u0435\u0435\u0442 \u043e\u0442\u043f\u0440\u0430\u0432\u043b\u044f\u0442\u044c \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f \u0432 \u0442\u043e\u043f\u0438\u043a. \u041e\u043d \u043f\u043e\u043d\u0430\u0434\u043e\u0431\u0438\u0442\u0441\u044f \u043d\u0430\u043c \u0432 \u0440\u0430\u043c\u043a\u0430\u0445 \u0442\u0435\u0441\u0442\u043e\u0432, \u043f\u043e\u044d\u0442\u043e\u043c\u0443 \u0434\u043e\u0431\u0430\u0432\u0438\u043c \u0437\u0430\u0432\u0438\u0441\u0438\u043c\u043e\u0441\u0442\u044c \u0432 \u0442\u0435\u0441\u0442\u0430\u0445:<\/p>\n<pre><code class=\"java\">testImplementation \"org.springframework.kafka:spring-kafka\"<\/code><\/pre>\n<p>\u042d\u0442\u043e\u0433\u043e \u043c\u043e\u0436\u043d\u043e \u0431\u044b\u043b\u043e \u0431\u044b \u0434\u043e\u0431\u0438\u0442\u044c\u0441\u044f \u0438 \u0430\u043b\u044c\u0442\u0435\u0440\u043d\u0430\u0442\u0438\u0432\u043d\u044b\u043c \u0441\u043f\u043e\u0441\u043e\u0431\u043e\u043c: \u0447\u0435\u0440\u0435\u0437 Spring-\u0438\u043d\u0438\u0446\u0438\u0430\u043b\u0438\u0437\u0430\u0442\u043e\u0440\u044b. \u041d\u0430\u043f\u0440\u0438\u043c\u0435\u0440, \u043c\u043e\u0436\u043d\u043e \u0431\u044b\u043b\u043e \u0431\u044b \u0441\u043e\u0437\u0434\u0430\u0442\u044c \u0441\u0432\u043e\u044e \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u044e \u0438\u043d\u0442\u0435\u0440\u0444\u0435\u0439\u0441\u0430 ApplicationContextInitializer. \u041d\u043e \u043c\u044b \u043d\u0435 \u0431\u0443\u0434\u0435\u043c \u0440\u0430\u0437\u0431\u0438\u0440\u0430\u0442\u044c \u044d\u0442\u043e\u0442 \u0441\u043f\u043e\u0441\u043e\u0431 \u0432 \u0441\u0442\u0430\u0442\u044c\u0435, \u043f\u043e\u0442\u043e\u043c\u0443 \u0447\u0442\u043e \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u044f \u0447\u0435\u0440\u0435\u0437 Extension \u0432\u044b\u0433\u043b\u044f\u0434\u0438\u0442 \u043a\u0440\u0430\u0441\u0438\u0432\u0435\u0435.<\/p>\n<p>\u0412 \u0438\u0442\u043e\u0433\u0435, \u0447\u0442\u043e\u0431\u044b \u043f\u0440\u0438\u043c\u0435\u043d\u044f\u0442\u044c \u0442\u0435\u043a\u0443\u0449\u0438\u0439 Extension \u0432 \u0442\u0435\u0441\u0442\u0430\u0445, \u043d\u0430\u043c \u0431\u044b\u043b\u0430 \u0431\u044b \u043f\u043e\u043b\u0435\u0437\u043d\u0430 \u0441\u0432\u043e\u044f \u0430\u043d\u043d\u043e\u0442\u0430\u0446\u0438\u044f \u043f\u043e \u0430\u043d\u0430\u043b\u043e\u0433\u0438\u0438 \u0441 \u0430\u043d\u043d\u043e\u0442\u0430\u0446\u0438\u0435\u0439 @WithFlinkCluster: <\/p>\n<pre><code class=\"java\">@Retention(RUNTIME) @ExtendWith({KafkaContainerExtension.class}) @Inherited public @interface WithKafkaContainer { }<\/code><\/pre>\n<p>\u042d\u0442\u0443 \u0430\u043d\u043d\u043e\u0442\u0430\u0446\u0438\u044e \u043c\u043e\u0436\u043d\u043e \u0432\u0435\u0448\u0430\u0442\u044c \u043d\u0430 \u043b\u044e\u0431\u043e\u0439 \u0442\u0435\u0441\u0442\u043e\u0432\u044b\u0439 \u043a\u043b\u0430\u0441\u0441, \u043a\u043e\u0442\u043e\u0440\u043e\u043c\u0443 \u043d\u0443\u0436\u043d\u0430 \u0438\u043d\u0442\u0435\u0433\u0440\u0430\u0446\u0438\u044f \u0441 Kafka.<\/p>\n<p><a class=\"anchor\" name=\"3\" id=\"3\"><\/a><\/p>\n<h2>\u0410\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u0438 \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u0434\u043b\u044f Kafka<\/h2>\n<p>\u0412\u043e \u0432\u0441\u0435\u0445 \u0442\u0435\u0441\u0442\u0430\u0445 \u0431\u044b\u043b\u043e \u0431\u044b \u0443\u0434\u043e\u0431\u043d\u043e \u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c\u0441\u044f \u0441\u0432\u043e\u0438\u043c\u0438 \u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u044f\u043c\u0438, \u0444\u0430\u0441\u0430\u0434\u0430\u043c\u0438 \u0438\u043b\u0438 dto \u0434\u043b\u044f \u0432\u0437\u0430\u0438\u043c\u043e\u0434\u0435\u0439\u0441\u0442\u0432\u0438\u044f \u0441 Kafka.<\/p>\n<p><a class=\"anchor\" name=\"4\" id=\"4\"><\/a><\/p>\n<h4>KafkaTestConsumer<\/h4>\n<p>\u0412\u043e-\u043f\u0435\u0440\u0432\u044b\u0445, \u043d\u0443\u0436\u043d\u043e \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0438\u0442\u044c \u0442\u0435\u0441\u0442\u043e\u0432\u044b\u0439 Consumer, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u0431\u0443\u0434\u0435\u043c \u0441\u043e\u0437\u0434\u0430\u0432\u0430\u0442\u044c \u0432 \u043a\u0430\u0436\u0434\u043e\u043c \u0442\u0435\u0441\u0442\u0435 \u043e\u0442\u0434\u0435\u043b\u044c\u043d\u043e, \u0447\u0442\u043e\u0431\u044b \u043e\u043d \u043f\u043e\u0434\u043a\u043b\u044e\u0447\u0430\u043b\u0441\u044f \u043a \u0442\u043e\u043f\u0438\u043a\u0430\u043c \u0432 \u0440\u0430\u043c\u043a\u0430\u0445 \u043d\u043e\u0432\u043e\u0439 \u043a\u043e\u043d\u0441\u044c\u044e\u043c\u0435\u0440\u043d\u043e\u0439 \u0433\u0440\u0443\u043f\u043f\u044b. \u0415\u0449\u0451 \u043e\u043d \u0434\u043e\u043b\u0436\u0435\u043d \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c \u0434\u0435-\/\u0441\u0435\u0440\u0438\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u044e Jackson, \u0432\u0435\u0434\u044c \u043c\u044b \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0438\u043b\u0438 \u0434\u043b\u044f \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0439 \u0444\u043e\u0440\u043c\u0430\u0442 JSON. \u0412\u0430\u0436\u043d\u043e \u043f\u043e\u043c\u043d\u0438\u0442\u044c: Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440 \u043f\u043e\u0434\u043d\u0438\u043c\u0430\u0435\u0442\u0441\u044f \u0432 \u0435\u0434\u0438\u043d\u0441\u0442\u0432\u0435\u043d\u043d\u043e\u043c \u044d\u043a\u0437\u0435\u043c\u043f\u043b\u044f\u0440\u0435. \u041f\u043e\u044d\u0442\u043e\u043c\u0443, \u0435\u0441\u043b\u0438 \u043d\u0430\u043f\u0438\u0441\u0430\u0442\u044c \u0442\u0435\u0441\u0442\u044b \u043d\u0435\u043f\u0440\u0430\u0432\u0438\u043b\u044c\u043d\u043e, \u043e\u043d\u0438 \u043c\u043e\u0433\u0443\u0442 \u043a\u043e\u0441\u0432\u0435\u043d\u043d\u043e \u0432\u043b\u0438\u044f\u0442\u044c \u0434\u0440\u0443\u0433 \u043d\u0430 \u0434\u0440\u0443\u0433\u0430. <\/p>\n<p>\u041a\u043b\u0430\u0441\u0441 \u0434\u043b\u044f \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u044f \u0442\u0435\u0441\u0442\u043e\u0432\u043e\u0433\u043e Kafka consumer \u043c\u043e\u0436\u043d\u043e \u043f\u0440\u0435\u0434\u0441\u0442\u0430\u0432\u0438\u0442\u044c \u0432 \u0442\u0430\u043a\u043e\u043c \u0432\u0438\u0434\u0435:<\/p>\n<pre><code class=\"java\">public class KafkaTestConsumer implements AutoCloseable {    private final Consumer&lt;String, String> consumer;    private final List&lt;KafkaMessage> receivedMessages = new CopyOnWriteArrayList&lt;>();    private final ObjectMapper objectMapper = createObjectMapper();     public KafkaTestConsumer(String bootstrapServers, Set&lt;String> topics) {        this.consumer = new KafkaConsumer&lt;>(            Map.of(                ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers,                ConsumerConfig.GROUP_ID_CONFIG, UUID.randomUUID().toString(),                ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName(),                ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName(),                ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 1,                ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, \"earliest\"            )        );        consumer.subscribe(topics);    }     public &lt;T> List&lt;T> receiveAndGetAll(String topic, Class&lt;T> clazz) {        return receiveAndGetAll()                   .stream()                   .filter(kafkaMessage -> topic.equals(kafkaMessage.getTopic()))                   .map(kafkaMessage -> readValue(kafkaMessage, clazz))                   .collect(toList());    }     private List&lt;KafkaMessage> receiveAndGetAll() {        final var records = consumer.poll(Duration.ofSeconds(5));        for (ConsumerRecord&lt;String, String> record : records) {            receivedMessages.add(new KafkaMessage(record.key(), record.topic(), record.value()));        }        consumer.commitSync();        return receivedMessages;    }     @SneakyThrows    private &lt;T> T readValue(KafkaMessage kafkaMessage, Class&lt;T> clazz) {        return objectMapper.readValue(kafkaMessage.getValue(), clazz);    }     @Override    public void close() {        receivedMessages.clear();        consumer.close();    } }<\/code><\/pre>\n<p>\u041a\u043e\u0434 \u0434\u043e\u0441\u0442\u0430\u0442\u043e\u0447\u043d\u043e \u043f\u0440\u043e\u0441\u0442. \u0412 \u043a\u043e\u043d\u0441\u0442\u0440\u0443\u043a\u0442\u043e\u0440\u0435 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u044f\u0435\u043c \u043e\u0441\u043d\u043e\u0432\u043d\u044b\u0435 \u043f\u0430\u0440\u0430\u043c\u0435\u0442\u0440\u044b \u043f\u043e\u0434\u043a\u043b\u044e\u0447\u0435\u043d\u0438\u044f \u043a Kafka, \u0441\u043e\u0437\u0434\u0430\u0451\u0442\u0441\u044f \u0431\u0430\u0437\u043e\u0432\u044b\u0439 consumer, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043f\u043e\u0434\u043f\u0438\u0441\u044b\u0432\u0430\u0435\u0442\u0441\u044f \u043d\u0430 \u043f\u0435\u0440\u0435\u0434\u0430\u043d\u043d\u044b\u0435 \u0442\u043e\u043f\u0438\u043a\u0438. \u041f\u043e\u0442\u043e\u043c \u043f\u0440\u0435\u0434\u043e\u0441\u0442\u0430\u0432\u043b\u044f\u0435\u043c \u043e\u0441\u043d\u043e\u0432\u043d\u043e\u0439 \u043c\u0435\u0442\u043e\u0434 receiveAndGetAll(), \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043e\u0442\u0434\u0430\u0451\u0442 \u0432\u0441\u0435 \u0434\u0435\u0441\u0435\u0440\u0438\u0430\u043b\u0438\u0437\u043e\u0432\u0430\u043d\u043d\u044b\u0435 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f \u0432 \u043a\u043e\u043d\u043a\u0440\u0435\u0442\u043d\u044b\u0439 \u0442\u0438\u043f \u0432 \u0440\u0430\u0437\u0440\u0435\u0437\u0435 \u043a\u0430\u043a\u043e\u0433\u043e-\u0442\u043e \u0442\u043e\u043f\u0438\u043a\u0430. \u0412\u043d\u0443\u0442\u0440\u0438 \u044d\u0442\u043e\u0433\u043e \u043c\u0435\u0442\u043e\u0434\u0430 \u043f\u0440\u043e\u0438\u0441\u0445\u043e\u0434\u0438\u0442 \u0432\u044b\u0437\u043e\u0432 \u0431\u0430\u0437\u043e\u0432\u043e\u0433\u043e consumer.poll() \u0434\u043b\u044f \u043e\u0431\u0440\u0430\u0449\u0435\u043d\u0438\u044f \u043a Kafka-\u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440\u0443. \u0412\u043d\u0443\u0442\u0440\u0438 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u043c \u0441\u0432\u043e\u044e dto KafkaMessage, \u0447\u0442\u043e\u0431\u044b \u043d\u0435 \u0440\u0430\u0431\u043e\u0442\u0430\u0442\u044c \u0441 ConsumerRecord \u043d\u0430\u043f\u0440\u044f\u043c\u0443\u044e, \u0432\u0435\u0434\u044c \u0431\u043e\u043b\u044c\u0448\u0438\u043d\u0441\u0442\u0432\u043e \u0438\u043d\u0444\u043e\u0440\u043c\u0430\u0446\u0438\u0438 \u0432 \u043d\u0451\u043c \u043d\u0430\u043c \u043d\u0435 \u0442\u0440\u0435\u0431\u0443\u0435\u0442\u0441\u044f:<\/p>\n<pre><code class=\"java\">@Value public class KafkaMessage {    String key;    String topic;    String value; }<\/code><\/pre>\n<p><a class=\"anchor\" name=\"5\" id=\"5\"><\/a><\/p>\n<h4>TestKafkaFacade<\/h4>\n<p>\u041d\u0430\u043c \u043f\u043e\u043d\u0430\u0434\u043e\u0431\u0438\u0442\u0441\u044f \u0430\u0431\u0441\u0442\u0440\u0430\u043a\u0446\u0438\u044f-\u0444\u0430\u0441\u0430\u0434,<\/p>\n<\/div>\n<\/div>\n<\/div>\n<\/div>\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-378097","post","type-post","status-publish","format-standard","hentry"],"_links":{"self":[{"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=\/wp\/v2\/posts\/378097","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=378097"}],"version-history":[{"count":0,"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=\/wp\/v2\/posts\/378097\/revisions"}],"wp:attachment":[{"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=%2Fwp%2Fv2%2Fmedia&parent=378097"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=%2Fwp%2Fv2%2Fcategories&post=378097"},{"taxonomy":"post_tag","embeddable":true,"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=%2Fwp%2Fv2%2Ftags&post=378097"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}