{"id":443829,"date":"2025-01-02T15:02:27","date_gmt":"2025-01-02T15:02:27","guid":{"rendered":"http:\/\/savepearlharbor.com\/?p=443829"},"modified":"-0001-11-30T00:00:00","modified_gmt":"2025-01-02T15:02:27","slug":"","status":"publish","type":"post","link":"https:\/\/savepearlharbor.com\/?p=443829","title":{"rendered":"<span>Kafka Streams \u04475: \u043d\u0430\u0441\u0442\u0440\u043e\u0439\u043a\u0430 \u043e\u043a\u043e\u043d, \u0440\u0430\u0431\u043e\u0442\u0430 \u0441 \u0437\u0430\u0434\u0435\u0440\u0436\u0430\u043d\u043d\u044b\u043c\u0438 \u0441\u043e\u0431\u044b\u0442\u0438\u044f\u043c\u0438 \u0438 suppression<\/span>"},"content":{"rendered":"<div><!--[--><!--]--><\/div>\n<div id=\"post-content-body\">\n<div>\n<div class=\"article-formatted-body article-formatted-body article-formatted-body_version-2\">\n<div xmlns=\"http:\/\/www.w3.org\/1999\/xhtml\">\n<p>\u042d\u0442\u043e \u043c\u043e\u044f \u0444\u0438\u043d\u0430\u043b\u044c\u043d\u0430\u044f \u0447\u0430\u0441\u0442\u044c(\u043d\u0443 \u043f\u043e\u043a\u0430 \u0447\u0442\u043e ;)) \u0441\u0435\u0440\u0438\u0438 \u0441\u0442\u0430\u0442\u0435\u0439 \u043f\u0440\u043e Kafka Streams, \u043f\u0440\u043e\u0448\u043b\u044b\u0435 \u0441\u0442\u0430\u0442\u044c\u0438 \u0442\u0443\u0442 [<a href=\"https:\/\/habr.com\/ru\/articles\/850832\/\" rel=\"noopener noreferrer nofollow\">\u043d\u043e\u043b\u044c<\/a>, <a href=\"https:\/\/habr.com\/ru\/articles\/854680\/\" rel=\"noopener noreferrer nofollow\">\u043e\u0434\u0438\u043d<\/a>, <a href=\"https:\/\/habr.com\/ru\/articles\/858668\/\" rel=\"noopener noreferrer nofollow\">\u0434\u0432\u0430<\/a>, <a href=\"https:\/\/habr.com\/ru\/articles\/862976\/\" rel=\"noopener noreferrer nofollow\">\u0442\u0440\u0438<\/a>] \u0422\u0435\u043f\u0435\u0440\u044c \u0434\u0430\u0432\u0430\u0439\u0442\u0435 \u0440\u0430\u0437\u0440\u0430\u0431\u043e\u0442\u0430\u0435\u043c \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0435, \u043a\u043e\u0442\u043e\u0440\u043e\u0435 \u0441\u0447\u0438\u0442\u044b\u0432\u0430\u0435\u0442:<\/p>\n<p>\u0421\u043e\u0431\u044b\u0442\u0438\u044f \u043e \u043f\u0443\u043b\u044c\u0441\u0435 \u0438\u0437 \u0442\u043e\u043f\u0438\u043a\u0430 <strong>pulse-events<\/strong>.<\/p>\n<pre><code class=\"json\">{ \"timestamp\": \"2024-12-26T09:02:00.000Z\" }<\/code><\/pre>\n<p>\u0421\u043e\u0431\u044b\u0442\u0438\u044f \u043e \u0442\u0435\u043c\u043f\u0435\u0440\u0430\u0442\u0443\u0440\u0435 \u0442\u0435\u043b\u0430 \u0438\u0437 \u0442\u043e\u043f\u0438\u043a\u0430 <strong>body-temp-events<\/strong>.<\/p>\n<pre><code class=\"json\">{ \"timestamp\": \"2024-12-26T09:02:00.000Z\", \"temperature\": 37.4, \"unit\": \"\u0421\" }<\/code><\/pre>\n<p>\u041f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0435 \u0431\u0443\u0434\u0435\u0442 \u043e\u0431\u0440\u0430\u0431\u0430\u0442\u044b\u0432\u0430\u0442\u044c \u044d\u0442\u0438 \u0441\u043e\u0431\u044b\u0442\u0438\u044f, \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u044f <strong>Kafka Streams<\/strong>. \u041c\u044b \u0431\u0443\u0434\u0435\u043c \u043e\u0440\u0438\u0435\u043d\u0442\u0438\u0440\u043e\u0432\u0430\u0442\u044c\u0441\u044f \u043d\u0430 <strong>event time<\/strong> (\u0432\u0440\u0435\u043c\u044f, \u0441\u043e\u0434\u0435\u0440\u0436\u0430\u0449\u0435\u0435\u0441\u044f \u0432 \u0441\u0430\u043c\u043e\u043c \u0441\u043e\u0431\u044b\u0442\u0438\u0438), \u0430 \u043d\u0435 \u043d\u0430 \u0432\u0440\u0435\u043c\u044f \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438 \u0438\u043b\u0438 \u0432\u0440\u0435\u043c\u044f \u0437\u0430\u0433\u0440\u0443\u0437\u043a\u0438 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f \u0432 Kafka.<\/p>\n<hr\/>\n<p><strong>\u041f\u043e\u043d\u044f\u0442\u0438\u044f \u0432\u0440\u0435\u043c\u0435\u043d\u0438 \u0432 Kafka Streams<\/strong><\/p>\n<p>\u0412 \u043f\u043e\u0442\u043e\u043a\u043e\u0432\u043e\u0439 \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0435 \u043e\u0431\u044b\u0447\u043d\u043e \u0432\u044b\u0434\u0435\u043b\u044f\u044e\u0442 \u0442\u0440\u0438 \u0432\u0438\u0434\u0430 \u0432\u0440\u0435\u043c\u0435\u043d\u0438:<\/p>\n<ul>\n<li>\n<p><strong>Event Time<\/strong> (\u0432\u0440\u0435\u043c\u044f \u0441\u043e\u0431\u044b\u0442\u0438\u044f). \u041a\u043e\u0433\u0434\u0430 \u0441\u043e\u0431\u044b\u0442\u0438\u0435 \u0440\u0435\u0430\u043b\u044c\u043d\u043e \u043f\u0440\u043e\u0438\u0437\u043e\u0448\u043b\u043e (\u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0435 \u043f\u043e\u043b\u044f timestamp, \u043f\u0440\u043e\u0441\u0442\u0430\u0432\u043b\u0435\u043d\u043d\u043e\u0435 \u0434\u0430\u0442\u0447\u0438\u043a\u043e\u043c \u0438\u043b\u0438 \u0441\u0435\u0440\u0432\u0438\u0441\u043e\u043c).<\/p>\n<\/li>\n<li>\n<p><strong>Ingestion Time<\/strong> (\u0438\u043b\u0438 Loading Time). \u0412\u0440\u0435\u043c\u044f, \u043a\u043e\u0433\u0434\u0430 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0435 \u043f\u043e\u043f\u0430\u0434\u0430\u0435\u0442 \u0432 \u0441\u0438\u0441\u0442\u0435\u043c\u0443 (Kafka), \u0444\u0438\u043a\u0441\u0438\u0440\u0443\u0435\u0442\u0441\u044f \u0431\u0440\u043e\u043a\u0435\u0440\u043e\u043c Kafka.<\/p>\n<\/li>\n<li>\n<p><strong>Processing Time<\/strong> (\u0432\u0440\u0435\u043c\u044f \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438). \u0422\u0435\u043a\u0443\u0449\u0435\u0435 \u0441\u0438\u0441\u0442\u0435\u043c\u043d\u043e\u0435 \u0432\u0440\u0435\u043c\u044f \u043f\u0440\u0438 \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0435 \u0441\u043e\u0431\u044b\u0442\u0438\u044f \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0435\u043c (JVM time).<\/p>\n<\/li>\n<\/ul>\n<p>\u0418\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u044f <strong>event time<\/strong>, \u0447\u0442\u043e\u0431\u044b \u0430\u043d\u0430\u043b\u0438\u0437\u0438\u0440\u043e\u0432\u0430\u0442\u044c \u0441\u043e\u0431\u044b\u0442\u0438\u044f \u0432 \u0442\u043e\u043c \u043f\u043e\u0440\u044f\u0434\u043a\u0435, \u0432 \u043a\u0430\u043a\u043e\u043c \u043e\u043d\u0438 \u0441\u043b\u0443\u0447\u0438\u043b\u0438\u0441\u044c \u0432 \u0440\u0435\u0430\u043b\u044c\u043d\u043e\u043c \u043c\u0438\u0440\u0435. \u0415\u0441\u043b\u0438 \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u044b\u0435 \u043c\u0435\u0442\u043a\u0438 \u043d\u0435\u0432\u0435\u0440\u043d\u044b\u0435 \u0438\u043b\u0438 \u043e\u0442\u0441\u0443\u0442\u0441\u0442\u0432\u0443\u044e\u0442, \u043c\u043e\u0436\u043d\u043e \u043e\u0440\u0438\u0435\u043d\u0442\u0438\u0440\u043e\u0432\u0430\u0442\u044c\u0441\u044f \u043d\u0430 \u0432\u0440\u0435\u043c\u044f \u0437\u0430\u0433\u0440\u0443\u0437\u043a\u0438 \u0438\u043b\u0438 \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438.<\/p>\n<hr\/>\n<p><strong>\u041f\u043e\u0434\u0433\u043e\u0442\u043e\u0432\u043a\u0430 \u043e\u043a\u0440\u0443\u0436\u0435\u043d\u0438\u044f \u0438 \u0441\u0442\u0440\u0443\u043a\u0442\u0443\u0440\u0430 \u043f\u0440\u043e\u0435\u043a\u0442\u0430<\/strong><\/p>\n<p><strong>\u041d\u0435\u043e\u0431\u0445\u043e\u0434\u0438\u043c\u044b\u0435 \u0438\u043d\u0441\u0442\u0440\u0443\u043c\u0435\u043d\u0442\u044b<\/strong><\/p>\n<ul>\n<li>\n<p><strong>Apache Kafka<\/strong> (\u0437\u0430\u043f\u0443\u0449\u0435\u043d\u043d\u0430\u044f \u043b\u043e\u043a\u0430\u043b\u044c\u043d\u043e \u0438\u043b\u0438 \u0432 Docker)<\/p>\n<\/li>\n<li>\n<p><strong>Gradle<\/strong> (\u043c\u044b \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u043c Kotlin DSL \u2014 build.gradle.kts)<\/p>\n<\/li>\n<li>\n<p><strong>Java 21<\/strong><\/p>\n<\/li>\n<\/ul>\n<p>\u0423\u0431\u0435\u0434\u0438\u0442\u0435\u0441\u044c, \u0447\u0442\u043e Kafka \u0437\u0430\u043f\u0443\u0449\u0435\u043d\u0430 \u043d\u0430 localhost:9092.<\/p>\n<p><strong>\u0421\u0442\u0440\u0443\u043a\u0442\u0443\u0440\u0430 \u043f\u0440\u043e\u0435\u043a\u0442\u0430<\/strong><\/p>\n<p><em>kafka-streams-medical-app<\/em><\/p>\n<p>\u251c\u2500 src<\/p>\n<p>\u2502\u00a0 \u251c\u2500 main<\/p>\n<p>\u2502\u00a0 \u2502\u00a0 \u251c\u2500 java<\/p>\n<p>\u2502\u00a0 \u2502\u00a0 \u2502\u00a0 \u2514\u2500 com.example.kafka<\/p>\n<p>\u2502\u00a0 \u2502\u00a0 \u2502 \u00a0 \u00a0 \u251c\u2500 MedicalMonitorApp.java<\/p>\n<p>\u2502\u00a0 \u2502\u00a0 \u2502 \u00a0 \u00a0 \u251c\u2500 PulseEvent.java<\/p>\n<p>\u2502\u00a0 \u2502\u00a0 \u2502 \u00a0 \u00a0 \u251c\u2500 BodyTempEvent.java<\/p>\n<p>\u2502\u00a0 \u2502\u00a0 \u2502 \u00a0 \u00a0 \u251c\u2500 MyTimestampExtractor.java<\/p>\n<p>\u2502\u00a0 \u2502\u00a0 \u2502 \u00a0 \u00a0 \u2514\u2500 &#8230;<\/p>\n<p>\u2502\u00a0 \u2502\u00a0 \u2514\u2500 resources<\/p>\n<p>\u2502\u00a0 \u2514\u2500 test<\/p>\n<p>\u2502 \u00a0 \u00a0 \u2514\u2500 java<\/p>\n<p>\u251c\u2500 build.gradle.kts<\/p>\n<p>\u2514\u2500 settings.gradle.kts<\/p>\n<p><strong>\u041f\u0440\u0438\u043c\u0435\u0440 build.gradle.kts<\/strong><\/p>\n<p>\u0423\u043f\u0440\u043e\u0449\u0451\u043d\u043d\u0430\u044f \u0432\u0435\u0440\u0441\u0438\u044f \u0441\u0431\u043e\u0440\u043e\u0447\u043d\u043e\u0433\u043e \u0444\u0430\u0439\u043b\u0430 (Kotlin DSL):<\/p>\n<pre><code class=\"java\">plugins { kotlin(\"jvm\") application  }  repositories { mavenCentral() }  dependencies { implementation(\"org.apache.kafka:kafka-streams:3.8.0\") implementation(\"com.fasterxml.jackson.core:jackson-databind:2.15.0\") implementation(\"com.fasterxml.jackson.module:jackson-module-kotlin:2.15.0\") implementation(\"org.slf4j:slf4j-api:2.0.0\")  runtimeOnly(\"org.slf4j:slf4j-simple:2.0.0\")  testImplementation(kotlin(\"test\")) testImplementation(\"org.apache.kafka:kafka-streams-test-utils:3.8.0\") }  application { mainClass.set(\"com.example.kafka.MedicalMonitorApp\") }<\/code><\/pre>\n<p>\u0417\u0430\u043f\u0443\u0441\u043a \u0441\u0431\u043e\u0440\u043a\u0438: <code>.\/gradlew clean build<\/code><\/p>\n<hr\/>\n<p><strong>\u0421\u043e\u0437\u0434\u0430\u043d\u0438\u0435 \u0438 \u043d\u0430\u0441\u0442\u0440\u043e\u0439\u043a\u0430 \u0442\u043e\u043f\u0438\u043a\u043e\u0432<\/strong><\/p>\n<p>\u0414\u043b\u044f \u0440\u0430\u0431\u043e\u0442\u044b \u043d\u0443\u0436\u043d\u043e \u0441\u043e\u0437\u0434\u0430\u0442\u044c \u0434\u0432\u0430 \u0442\u043e\u043f\u0438\u043a\u0430:<\/p>\n<ul>\n<li>\n<p><strong>pulse-events<\/strong><\/p>\n<\/li>\n<li>\n<p><strong>body-temp-events<\/strong><\/p>\n<\/li>\n<\/ul>\n<p>\u041f\u0440\u0438\u043c\u0435\u0440 \u043a\u043e\u043c\u0430\u043d\u0434 (\u043a\u043e\u043d\u0441\u043e\u043b\u044c Kafka):<\/p>\n<pre><code class=\"bash\">bin\/kafka-topics.sh --bootstrap-server localhost:9092 \\\\ --create --topic pulse-events \\\\ --partitions 1 --replication-factor 1  bin\/kafka-topics.sh --bootstrap-server localhost:9092 \\\\ --create --topic body-temp-events \\\\ --partitions 1 --replication-factor 1<\/code><\/pre>\n<p><strong>\u0421\u043e\u0437\u0434\u0430\u043d\u0438\u0435 \u043a\u043b\u0430\u0441\u0441\u043e\u0432-\u0441\u0443\u0449\u043d\u043e\u0441\u0442\u0435\u0439 (DTO)<\/strong><\/p>\n<p>\u041f\u0440\u0438\u043c\u0435\u0440 \u043a\u043b\u0430\u0441\u0441\u0430 <strong>PulseEvent.java<\/strong>:<\/p>\n<pre><code class=\"java\">public record PulseEvent(String timestamp) { }<\/code><\/pre>\n<p>\u041f\u0440\u0438\u043c\u0435\u0440 \u043a\u043b\u0430\u0441\u0441\u0430 <strong>BodyTempEvent.java<\/strong>:<\/p>\n<pre><code class=\"java\">public record BodyTempEvent(String timestamp, double temperature, String unit) { }<\/code><\/pre>\n<hr\/>\n<p><strong>\u041a\u0430\u043a \u0440\u0430\u0437\u043b\u0438\u0447\u0430\u0442\u044c \u0432\u0440\u0435\u043c\u044f \u0441\u043e\u0431\u044b\u0442\u0438\u044f, \u0432\u0440\u0435\u043c\u044f \u0437\u0430\u0433\u0440\u0443\u0437\u043a\u0438 \u0438 \u0432\u0440\u0435\u043c\u044f \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438<\/strong><\/p>\n<p>\u041c\u044b \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u043c <strong>event time<\/strong> (\u0438\u0437 \u043f\u043e\u043b\u044f timestamp \u0432\u043e \u0432\u0445\u043e\u0434\u044f\u0449\u0435\u043c \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0438), \u0447\u0442\u043e\u0431\u044b \u0430\u043d\u0430\u043b\u0438\u0437\u0438\u0440\u043e\u0432\u0430\u0442\u044c \u0441\u043e\u0431\u044b\u0442\u0438\u044f \u0432 \u0440\u0435\u0430\u043b\u044c\u043d\u043e\u043c \u043f\u043e\u0440\u044f\u0434\u043a\u0435 \u0438\u0445 \u0432\u043e\u0437\u043d\u0438\u043a\u043d\u043e\u0432\u0435\u043d\u0438\u044f. Kafka \u043f\u043e \u0443\u043c\u043e\u043b\u0447\u0430\u043d\u0438\u044e \u043c\u043e\u0436\u0435\u0442 \u0431\u0440\u0430\u0442\u044c \u0432\u0440\u0435\u043c\u044f \u0437\u0430\u0433\u0440\u0443\u0437\u043a\u0438 (ingestion), \u043d\u043e \u0434\u043b\u044f \u043c\u0435\u0434\u0438\u0446\u0438\u043d\u0441\u043a\u0438\u0445 \u043f\u043e\u043a\u0430\u0437\u0430\u0442\u0435\u043b\u0435\u0439 \u0432\u0430\u0436\u043d\u043e \u0443\u0447\u0438\u0442\u044b\u0432\u0430\u0442\u044c \u0444\u0430\u043a\u0442\u0438\u0447\u0435\u0441\u043a\u043e\u0435 \u0432\u0440\u0435\u043c\u044f \u0438\u0437\u043c\u0435\u0440\u0435\u043d\u0438\u044f.<\/p>\n<p>\u0415\u0441\u043b\u0438 \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u044b\u0435 \u043c\u0435\u0442\u043a\u0438 \u043d\u0435\u043a\u043e\u0440\u0440\u0435\u043a\u0442\u043d\u044b (\u0438\u043b\u0438 \u043e\u0442\u0441\u0443\u0442\u0441\u0442\u0432\u0443\u044e\u0442), \u043c\u043e\u0436\u043d\u043e \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c fallback \u043d\u0430 ingestion time \u0438\u043b\u0438 \u043d\u0430 processing time (\u0441\u0438\u0441\u0442\u0435\u043c\u043d\u043e\u0435 \u0432\u0440\u0435\u043c\u044f \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438).<\/p>\n<hr\/>\n<p><strong>\u0421\u043e\u0437\u0434\u0430\u043d\u0438\u0435 \u043a\u0430\u0441\u0442\u043e\u043c\u043d\u043e\u0433\u043e TimestampExtractor<\/strong><\/p>\n<p>\u0427\u0442\u043e\u0431\u044b Kafka Streams \u0437\u043d\u0430\u043b, \u0447\u0442\u043e \u043c\u044b \u043e\u0440\u0438\u0435\u043d\u0442\u0438\u0440\u0443\u0435\u043c\u0441\u044f \u043d\u0430 <strong>event time<\/strong>, \u043d\u0443\u0436\u043d\u043e \u043d\u0430\u0441\u0442\u0440\u043e\u0438\u0442\u044c TimestampExtractor, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u0431\u0443\u0434\u0435\u0442 \u0437\u0430\u0431\u0438\u0440\u0430\u0442\u044c \u043f\u043e\u043b\u0435 timestamp \u0438\u0437 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f:<\/p>\n<pre><code class=\"java\">import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.streams.processor.TimestampExtractor; import com.fasterxml.jackson.databind.ObjectMapper;  public class MyTimestampExtractor implements TimestampExtractor {      private final ObjectMapper objectMapper = new ObjectMapper();      @Override     public long extract(final ConsumerRecord&lt;Object, Object&gt; record, final long partitionTime) {         if (record.value() == null) {             return partitionTime;         }          try {             String jsonString = record.value().toString();             String timestampStr = objectMapper.readTree(jsonString).get(\"timestamp\").asText();             return javax.xml.bind.DatatypeConverter.parseDateTime(timestampStr).getTimeInMillis();         } catch (Exception e) {             return partitionTime;         }     } }<\/code><\/pre>\n<ul>\n<li>\n<p>\u0418\u0437\u0432\u043b\u0435\u043a\u0430\u0435\u043c JSON-\u043f\u043e\u043b\u0435 &#171;timestamp&#187; \u0438\u0437 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f.<\/p>\n<\/li>\n<li>\n<p>\u041f\u0440\u0435\u043e\u0431\u0440\u0430\u0437\u0443\u0435\u043c \u0435\u0433\u043e \u0438\u0437 ISO-8601 \u0432 \u043c\u0438\u043b\u043b\u0438\u0441\u0435\u043a\u0443\u043d\u0434\u044b (long).<\/p>\n<\/li>\n<li>\n<p>\u0415\u0441\u043b\u0438 \u043f\u0430\u0440\u0441\u0438\u043d\u0433 \u043d\u0435 \u0443\u0434\u0430\u043b\u0441\u044f, \u0432\u043e\u0437\u0432\u0440\u0430\u0449\u0430\u0435\u043c partitionTime.<\/p>\n<\/li>\n<\/ul>\n<hr\/>\n<p><strong>\u041a\u0430\u043a \u0432\u0440\u0435\u043c\u044f \u0443\u043f\u0440\u0430\u0432\u043b\u044f\u0435\u0442 \u043f\u043e\u0442\u043e\u043a\u043e\u043c \u0434\u0430\u043d\u043d\u044b\u0445 \u0432 Kafka Streams<\/strong><\/p>\n<p>\u0411\u043b\u0430\u0433\u043e\u0434\u0430\u0440\u044f \u043a\u0430\u0441\u0442\u043e\u043c\u043d\u043e\u043c\u0443 TimestampExtractor, Kafka Streams \u0432\u044b\u0441\u0442\u0440\u0430\u0438\u0432\u0430\u0435\u0442 \u0441\u043e\u0431\u044b\u0442\u0438\u044f \u043f\u043e \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u043e\u0439 \u0448\u043a\u0430\u043b\u0435 \u043d\u0430 \u043e\u0441\u043d\u043e\u0432\u0435 <strong>event time<\/strong>. \u0415\u0441\u043b\u0438 \u043f\u0440\u0438\u0445\u043e\u0434\u044f\u0442 \u201c\u0437\u0430\u043f\u0430\u0437\u0434\u044b\u0432\u0430\u044e\u0449\u0438\u0435\u201d \u0441\u043e\u0431\u044b\u0442\u0438\u044f \u0441\u043e \u0441\u0442\u0430\u0440\u044b\u043c\u0438 \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u044b\u043c\u0438 \u043c\u0435\u0442\u043a\u0430\u043c\u0438, \u043e\u043d\u0438 \u043c\u043e\u0433\u0443\u0442 \u043f\u043e\u043f\u0430\u0441\u0442\u044c \u0432 \u0443\u0436\u0435 \u0437\u0430\u043a\u0440\u044b\u0442\u044b\u0435 \u043e\u043a\u043d\u0430 (\u044d\u0442\u043e \u0440\u0435\u0433\u0443\u043b\u0438\u0440\u0443\u0435\u0442\u0441\u044f \u043d\u0430\u0441\u0442\u0440\u043e\u0439\u043a\u0430\u043c\u0438 grace period).<\/p>\n<p>\u042d\u0442\u043e \u043e\u0441\u043e\u0431\u0435\u043d\u043d\u043e \u0432\u0430\u0436\u043d\u043e, \u043a\u043e\u0433\u0434\u0430 \u0440\u0435\u0430\u043b\u044c\u043d\u043e\u0435 \u043f\u043e\u043b\u0443\u0447\u0435\u043d\u0438\u0435 \u0441\u043e\u0431\u044b\u0442\u0438\u044f \u0437\u0430\u0434\u0435\u0440\u0436\u0438\u0432\u0430\u0435\u0442\u0441\u044f, \u043d\u043e \u043c\u044b \u0445\u043e\u0442\u0438\u043c \u0443\u0447\u0438\u0442\u044b\u0432\u0430\u0442\u044c \u0435\u0433\u043e \u0438\u043c\u0435\u043d\u043d\u043e \u0432 \u0442\u043e\u043c \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u043e\u043c \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u0435, \u043a \u043a\u043e\u0442\u043e\u0440\u043e\u043c\u0443 \u043e\u043d\u043e \u043f\u0440\u0438\u043d\u0430\u0434\u043b\u0435\u0436\u0438\u0442 \u0444\u0430\u043a\u0442\u0438\u0447\u0435\u0441\u043a\u0438.<\/p>\n<hr\/>\n<p><strong>\u0422\u0438\u043f\u044b \u043e\u043a\u043e\u043d \u0432 Kafka Streams<\/strong><\/p>\n<p><strong>Tumbling Windows<\/strong><\/p>\n<ul>\n<li>\n<p>\u041e\u043a\u043d\u0430 \u0444\u0438\u043a\u0441\u0438\u0440\u043e\u0432\u0430\u043d\u043d\u043e\u0433\u043e \u0440\u0430\u0437\u043c\u0435\u0440\u0430, \u043a\u043e\u0442\u043e\u0440\u044b\u0435 <strong>\u043d\u0435 \u043f\u0435\u0440\u0435\u043a\u0440\u044b\u0432\u0430\u044e\u0442\u0441\u044f<\/strong>.<\/p>\n<\/li>\n<li>\n<p>\u041d\u0430\u043f\u0440\u0438\u043c\u0435\u0440, \u043e\u043a\u043d\u043e \u0432 1 \u043c\u0438\u043d\u0443\u0442\u0443: [09:00 \u2013 09:01], [09:01 \u2013 09:02] \u0438 \u0442.\u0434.<\/p>\n<\/li>\n<\/ul>\n<p><strong>Hopping Windows<\/strong><\/p>\n<ul>\n<li>\n<p>\u041e\u043a\u043d\u0430 \u0444\u0438\u043a\u0441\u0438\u0440\u043e\u0432\u0430\u043d\u043d\u043e\u0433\u043e \u0440\u0430\u0437\u043c\u0435\u0440\u0430, <strong>\u043d\u043e \u043f\u0435\u0440\u0435\u043a\u0440\u044b\u0432\u0430\u044e\u0449\u0438\u0435\u0441\u044f<\/strong>.<\/p>\n<\/li>\n<li>\n<p>\u041f\u0440\u0438\u043c\u0435\u0440: \u0440\u0430\u0437\u043c\u0435\u0440 \u043e\u043a\u043d\u0430 1 \u043c\u0438\u043d\u0443\u0442\u0430, \u0441\u0434\u0432\u0438\u0433 30 \u0441\u0435\u043a\u0443\u043d\u0434, \u043f\u043e\u043b\u0443\u0447\u0430\u0435\u043c [09:00 \u2013 09:01], [09:00:30 \u2013 09:01:30] \u0438 \u0442.\u0434.<\/p>\n<\/li>\n<\/ul>\n<p><strong>Session Windows<\/strong><\/p>\n<ul>\n<li>\n<p>\u041e\u043a\u043d\u0430, \u043a\u043e\u0442\u043e\u0440\u044b\u0435 <strong>\u043d\u0435 \u0438\u043c\u0435\u044e\u0442 \u0444\u0438\u043a\u0441\u0438\u0440\u043e\u0432\u0430\u043d\u043d\u043e\u0439 \u0434\u043b\u0438\u0442\u0435\u043b\u044c\u043d\u043e\u0441\u0442\u0438<\/strong>; \u043e\u043d\u0438 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u044f\u044e\u0442\u0441\u044f \u201c\u043f\u0430\u0443\u0437\u0430\u043c\u0438\u201d \u043c\u0435\u0436\u0434\u0443 \u0441\u043e\u0431\u044b\u0442\u0438\u044f\u043c\u0438.<\/p>\n<\/li>\n<li>\n<p>\u0415\u0441\u043b\u0438 \u043c\u0435\u0436\u0434\u0443 \u0441\u043e\u0431\u044b\u0442\u0438\u044f\u043c\u0438 \u043d\u0435\u0442 \u0431\u043e\u043b\u044c\u0448\u0438\u0445 \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u043e\u0432, \u043e\u043a\u043d\u043e \u0440\u0430\u0441\u0448\u0438\u0440\u044f\u0435\u0442\u0441\u044f; \u043f\u0440\u0438 \u0434\u043b\u0438\u0442\u0435\u043b\u044c\u043d\u043e\u043c \u0437\u0430\u0442\u0438\u0448\u044c\u0435 \u043e\u043a\u043d\u043e \u201c\u0437\u0430\u043a\u0440\u044b\u0432\u0430\u0435\u0442\u0441\u044f\u201d.<\/p>\n<\/li>\n<\/ul>\n<p><strong>Sliding Join Windows<\/strong><\/p>\n<ul>\n<li>\n<p>\u041e\u043a\u043d\u0430, \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u043c\u044b\u0435 \u0432 \u043e\u043f\u0435\u0440\u0430\u0446\u0438\u044f\u0445 join \u043f\u043e\u0442\u043e\u043a\u043e\u0432, \u0433\u0434\u0435 \u043d\u0430\u043c \u0432\u0430\u0436\u043d\u043e \u201c\u0437\u0430\u0445\u0432\u0430\u0442\u0438\u0442\u044c\u201d \u0441\u043e\u0431\u044b\u0442\u0438\u044f, \u043a\u043e\u0442\u043e\u0440\u044b\u0435 \u043f\u0440\u043e\u0438\u0437\u043e\u0448\u043b\u0438 \u0432 \u043f\u0440\u043e\u043c\u0435\u0436\u0443\u0442\u043a\u0435 [t, t+N].<\/p>\n<\/li>\n<li>\n<p>\u041f\u043e\u0445\u043e\u0436\u0435 \u043d\u0430 hopping windows, \u043d\u043e \u043f\u0440\u0438\u043c\u0435\u043d\u044f\u0435\u0442\u0441\u044f \u0434\u043b\u044f Stream-Stream join.<\/p>\n<\/li>\n<\/ul>\n<p><strong>\u0412\u044b\u043f\u043e\u043b\u043d\u0435\u043d\u0438\u0435 \u0438 \u043f\u0440\u043e\u0438\u0437\u0432\u043e\u0434\u0441\u0442\u0432\u043e \u043e\u043a\u043e\u043d\u043d\u044b\u0445 \u0441\u043e\u0435\u0434\u0438\u043d\u0435\u043d\u0438\u0439 (Stream-Stream join)<\/strong><\/p>\n<p>\u0427\u0442\u043e\u0431\u044b \u043e\u0431\u044a\u0435\u0434\u0438\u043d\u0438\u0442\u044c \u0434\u0432\u0430 \u043f\u043e\u0442\u043e\u043a\u0430 \u043d\u0430 \u043e\u0441\u043d\u043e\u0432\u0435 \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u044b\u0445 \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u043e\u0432, \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u0442\u0441\u044f \u043e\u043a\u043e\u043d\u043d\u043e\u0435 \u0441\u043e\u0435\u0434\u0438\u043d\u0435\u043d\u0438\u0435 (Stream-Stream Join). \u041d\u0430\u043f\u0440\u0438\u043c\u0435\u0440, \u0445\u043e\u0442\u0438\u043c \u201c\u0441\u043a\u043b\u0435\u0438\u0442\u044c\u201d \u0434\u0430\u043d\u043d\u044b\u0435 \u043e \u043f\u0443\u043b\u044c\u0441\u0435 \u0438 \u0442\u0435\u043c\u043f\u0435\u0440\u0430\u0442\u0443\u0440\u0435, \u043f\u043e\u0441\u0442\u0443\u043f\u0438\u0432\u0448\u0438\u0435 \u0432 \u043f\u0440\u0435\u0434\u0435\u043b\u0430\u0445 \u043e\u0434\u043d\u043e\u0439 \u043c\u0438\u043d\u0443\u0442\u044b:<\/p>\n<pre><code class=\"java\">KStream&lt;String, PulseEvent&gt; pulseStream = builder     .stream(\"pulse-events\", Consumed.with(Serdes.String(), pulseSerde)         .withTimestampExtractor(new MyTimestampExtractor()));  KStream&lt;String, BodyTempEvent&gt; tempStream = builder     .stream(\"body-temp-events\", Consumed.with(Serdes.String(), bodyTempSerde)         .withTimestampExtractor(new MyTimestampExtractor()));  KStream&lt;String, CombinedMeasurement&gt; joinedStream = pulseStream.join(     tempStream,     (pulseValue, tempValue) -&gt; {         CombinedMeasurement cm = new CombinedMeasurement();         cm.setTimestampPulse(pulseValue.getTimestamp());         cm.setTimestampTemp(tempValue.getTimestamp());         cm.setTemperature(tempValue.getTemperature());         return cm;     },     JoinWindows.ofTimeDifferenceWithNoGrace(Duration.ofMinutes(1)),     StreamJoined.with(Serdes.String(), pulseSerde, bodyTempSerde) );<\/code><\/pre>\n<p>\u0422\u0430\u043a\u0438\u043c \u043e\u0431\u0440\u0430\u0437\u043e\u043c, \u0435\u0441\u043b\u0438 \u0441\u043e\u0431\u044b\u0442\u0438\u044f \u043f\u043e \u043a\u043b\u044e\u0447\u0443 \u0441\u043e\u0432\u043f\u0430\u0434\u0430\u044e\u0442 \u0438 \u043f\u043e\u043f\u0430\u0434\u0430\u044e\u0442 \u0432 \u043e\u0431\u0449\u0438\u0439 \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u043e\u0439 \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b, \u043e\u043d\u0438 \u0431\u0443\u0434\u0443\u0442 \u043e\u0431\u044a\u0435\u0434\u0438\u043d\u0435\u043d\u044b \u0432 \u0435\u0434\u0438\u043d\u044b\u0439 \u0440\u0435\u0437\u0443\u043b\u044c\u0442\u0430\u0442.<\/p>\n<hr\/>\n<p><strong>\u0417\u0430\u0447\u0435\u043c \u043d\u0443\u0436\u0435\u043d \u043e\u043f\u0435\u0440\u0430\u0442\u043e\u0440 suppress \u0438 \u043a\u0430\u043a \u0435\u0433\u043e \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c<\/strong><\/p>\n<p>\u0412 KTable \u0438 \u0430\u0433\u0440\u0435\u0433\u0438\u0440\u0443\u044e\u0449\u0438\u0445 \u043e\u043f\u0435\u0440\u0430\u0446\u0438\u044f\u0445 Kafka Streams \u0432\u043e\u0437\u043d\u0438\u043a\u0430\u044e\u0442 \u043f\u0440\u043e\u043c\u0435\u0436\u0443\u0442\u043e\u0447\u043d\u044b\u0435 \u0440\u0435\u0437\u0443\u043b\u044c\u0442\u0430\u0442\u044b (updates). \u041e\u043f\u0435\u0440\u0430\u0442\u043e\u0440 suppress \u043f\u043e\u0437\u0432\u043e\u043b\u044f\u0435\u0442 \u201c\u0437\u0430\u0434\u0435\u0440\u0436\u0438\u0432\u0430\u0442\u044c\u201d (\u043f\u043e\u0434\u0430\u0432\u043b\u044f\u0442\u044c) \u043f\u0443\u0431\u043b\u0438\u043a\u0430\u0446\u0438\u044e \u043f\u0440\u043e\u043c\u0435\u0436\u0443\u0442\u043e\u0447\u043d\u044b\u0445 \u0440\u0435\u0437\u0443\u043b\u044c\u0442\u0430\u0442\u043e\u0432 \u0438 \u0432\u044b\u0434\u0430\u0432\u0430\u0442\u044c <strong>\u0442\u043e\u043b\u044c\u043a\u043e \u043a\u043e\u043d\u0435\u0447\u043d\u044b\u0439<\/strong> (\u043f\u043e\u0441\u043b\u0435 \u0437\u0430\u043a\u0440\u044b\u0442\u0438\u044f \u043e\u043a\u043d\u0430).<\/p>\n<pre><code class=\"java\">KTable&lt;Windowed&lt;String&gt;, Long&gt; aggregatedTable = pulseStream     .groupByKey()     .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(1)))     .count();  KTable&lt;Windowed&lt;String&gt;, Long&gt; suppressedTable = aggregatedTable     .suppress(Suppressed.untilWindowCloses(Suppressed.BufferConfig.unbounded())); <\/code><\/pre>\n<ul>\n<li>\n<p>untilWindowCloses \u043e\u0437\u043d\u0430\u0447\u0430\u0435\u0442, \u0447\u0442\u043e \u043c\u044b \u043f\u0443\u0431\u043b\u0438\u043a\u0443\u0435\u043c \u0440\u0435\u0437\u0443\u043b\u044c\u0442\u0430\u0442 \u0442\u043e\u043b\u044c\u043a\u043e \u043e\u0434\u0438\u043d \u0440\u0430\u0437 \u2014 \u043a\u043e\u0433\u0434\u0430 \u043e\u043a\u043d\u043e \u043e\u043a\u043e\u043d\u0447\u0430\u0442\u0435\u043b\u044c\u043d\u043e \u0437\u0430\u043a\u0440\u044b\u0442\u043e.<\/p>\n<\/li>\n<li>\n<p>\u042d\u0442\u043e \u0438\u0437\u0431\u0430\u0432\u043b\u044f\u0435\u0442 \u043e\u0442 \u201c\u0448\u0442\u043e\u0440\u043c\u0430\u201d \u043f\u043e\u0441\u0442\u043e\u044f\u043d\u043d\u044b\u0445 \u043e\u0431\u043d\u043e\u0432\u043b\u0435\u043d\u0438\u0439.<\/p>\n<\/li>\n<\/ul>\n<hr\/>\n<p><strong>\u0427\u0442\u043e \u0432\u044b\u0448\u043b\u043e \u0432 \u0438\u0442\u043e\u0433\u0435<\/strong><\/p>\n<p>\u041d\u0438\u0436\u0435 \u0443\u043f\u0440\u043e\u0449\u0451\u043d\u043d\u044b\u0439 \u043a\u043e\u0434 <strong>MedicalMonitorApp.java<\/strong>, \u0433\u0434\u0435:<\/p>\n<ul>\n<li>\n<p>\u041d\u0430\u0441\u0442\u0440\u0430\u0438\u0432\u0430\u0435\u043c \u043a\u043e\u043d\u0444\u0438\u0433\u0443\u0440\u0430\u0446\u0438\u044e Kafka Streams.<\/p>\n<\/li>\n<li>\n<p>\u0421\u043e\u0437\u0434\u0430\u0451\u043c \u043f\u043e\u0442\u043e\u043a\u0438 \u0434\u043b\u044f \u043f\u0443\u043b\u044c\u0441\u0430 \u0438 \u0442\u0435\u043c\u043f\u0435\u0440\u0430\u0442\u0443\u0440\u044b.<\/p>\n<\/li>\n<li>\n<p>\u0412\u044b\u043f\u043e\u043b\u043d\u044f\u0435\u043c join \u043f\u043e \u043e\u043a\u043d\u0443 \u0432 1 \u043c\u0438\u043d\u0443\u0442\u0443.<\/p>\n<\/li>\n<li>\n<p>\u041e\u0442\u043f\u0440\u0430\u0432\u043b\u044f\u0435\u043c \u0440\u0435\u0437\u0443\u043b\u044c\u0442\u0430\u0442 \u0432 \u043d\u043e\u0432\u044b\u0439 \u0442\u043e\u043f\u0438\u043a.<\/p>\n<\/li>\n<\/ul>\n<pre><code class=\"java\">import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.*; import org.apache.kafka.streams.kstream.*; import java.time.Duration; import java.util.Properties;  public class MedicalMonitorApp {      public static void main(String[] args) {         Properties props = new Properties();         props.put(StreamsConfig.APPLICATION_ID_CONFIG, \"medical-monitor-app\");         props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, \"localhost:9092\");         props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.StringSerde.class);         props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.StringSerde.class);         props.put(StreamsConfig.DEFAULT_TIMESTAMP_EXTRACTOR_CLASS_CONFIG, MyTimestampExtractor.class);          StreamsBuilder builder = new StreamsBuilder();          KStream&lt;String, PulseEvent&gt; pulseStream = builder.stream(             \"pulse-events\",             Consumed.with(                 Serdes.String(),                 Serdes.serdeFrom(new JsonSerializer&lt;&gt;(), new JsonDeserializer&lt;&gt;(PulseEvent.class))             )         );          KStream&lt;String, BodyTempEvent&gt; tempStream = builder.stream(             \"body-temp-events\",             Consumed.with(                 Serdes.String(),                 Serdes.serdeFrom(new JsonSerializer&lt;&gt;(), new JsonDeserializer&lt;&gt;(BodyTempEvent.class))             )         );          KStream&lt;String, String&gt; joinedStream = pulseStream.join(             tempStream,             (pulseVal, tempVal) -&gt; {                 return \"Pulse timestamp: \" + pulseVal.getTimestamp() +                        \", Temp timestamp: \" + tempVal.getTimestamp() +                        \", Temp: \" + tempVal.getTemperature();             },             JoinWindows.ofTimeDifferenceWithNoGrace(Duration.ofMinutes(1)),             StreamJoined.with(                 Serdes.String(),                 Serdes.serdeFrom(new JsonSerializer&lt;&gt;(), new JsonDeserializer&lt;&gt;(PulseEvent.class)),                 Serdes.serdeFrom(new JsonSerializer&lt;&gt;(), new JsonDeserializer&lt;&gt;(BodyTempEvent.class))             )         );          joinedStream.to(\"joined-medical-events\", Produced.with(Serdes.String(), Serdes.String()));          KafkaStreams streams = new KafkaStreams(builder.build(), props);         Runtime.getRuntime().addShutdownHook(new Thread(streams::close));         streams.start();     } }<\/code><\/pre>\n<p>\u0412 \u043f\u0440\u0438\u043c\u0435\u0440\u0435 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u043d\u044b JsonSerializer \u0438 JsonDeserializer (\u043e\u0431\u0451\u0440\u0442\u043a\u0438 \u043d\u0430\u0434 Jackson), \u043a\u043e\u0442\u043e\u0440\u044b\u0435 \u043c\u043e\u0436\u043d\u043e \u0440\u0435\u0430\u043b\u0438\u0437\u043e\u0432\u0430\u0442\u044c \u0432\u0440\u0443\u0447\u043d\u0443\u044e \u0438\u043b\u0438 \u0432\u0437\u044f\u0442\u044c \u0433\u043e\u0442\u043e\u0432\u0443\u044e \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u044e.<\/p>\n<hr\/>\n<p><strong>\u0418\u043d\u0441\u0442\u0440\u0443\u043a\u0446\u0438\u0438 \u043f\u043e \u0437\u0430\u043f\u0443\u0441\u043a\u0443 \u0438 \u0442\u0435\u0441\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044e \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u044f<\/strong><\/p>\n<p><strong>\u0421\u0431\u043e\u0440\u043a\u0430<\/strong>: <code>.\/gradlew clean build<\/code><\/p>\n<p><strong>\u0417\u0430\u043f\u0443\u0441\u043a<\/strong> (\u0447\u0435\u0440\u0435\u0437 Gradle): <code>.\/gradlew run<\/code><\/p>\n<p><strong>\u041e\u0442\u043f\u0440\u0430\u0432\u043a\u0430 \u0442\u0435\u0441\u0442\u043e\u0432\u044b\u0445 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0439<\/strong>:<\/p>\n<pre><code class=\"bash\">bin\/kafka-console-producer.sh --broker-list localhost:9092 --topic pulse-events &gt; {\"timestamp\": \"2024-12-26T09:02:00.000Z\"} &gt; {\"timestamp\": \"2024-12-26T09:02:30.000Z\"}  bin\/kafka-console-producer.sh --broker-list localhost:9092 --topic body-temp-events &gt; {\"timestamp\": \"2024-12-26T09:02:05.000Z\", \"temperature\": 37.4, \"unit\": \"C\"} &gt; {\"timestamp\": \"2024-12-26T09:02:25.000Z\", \"temperature\": 37.1, \"unit\": \"C\"}<\/code><\/pre>\n<p><strong>\u041f\u0440\u043e\u0432\u0435\u0440\u043a\u0430 \u0440\u0435\u0437\u0443\u043b\u044c\u0442\u0430\u0442\u0430<\/strong> \u0432 \u0442\u043e\u043f\u0438\u043a\u0435 joined-medical-events:<\/p>\n<pre><code class=\"bash\">bin\/kafka-console-consumer.sh --bootstrap-server localhost:9092 \\\\ --topic joined-medical-events --from-beginning<\/code><\/pre>\n<p>\u041e\u0436\u0438\u0434\u0430\u0435\u043c\u044b\u0439 \u0432\u044b\u0432\u043e\u0434:<\/p>\n<pre><code>Pulse timestamp: 2024-12-26T09:02:00.000Z, Temp timestamp: 2024-12-26T09:02:05.000Z, Temp: 37.4 Pulse timestamp: 2024-12-26T09:02:30.000Z, Temp timestamp: 2024-12-26T09:02:25.000Z, Temp: 37.1<\/code><\/pre>\n<p>\u0421\u0442\u0430\u043b\u043e \u043f\u043e\u043d\u044f\u0442\u043d\u0435\u0435, \u0447\u0435\u043c \u043e\u0442\u043b\u0438\u0447\u0430\u044e\u0442\u0441\u044f event time, ingestion time \u0438 processing time, \u0438 \u043f\u043e\u0447\u0435\u043c\u0443 \u0442\u0430\u043a \u0432\u0430\u0436\u043d\u043e \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c \u0444\u0430\u043a\u0442\u0438\u0447\u0435\u0441\u043a\u043e\u0435 \u0432\u0440\u0435\u043c\u044f \u0432\u043e\u0437\u043d\u0438\u043a\u043d\u043e\u0432\u0435\u043d\u0438\u044f \u0441\u043e\u0431\u044b\u0442\u0438\u044f. \u041f\u043e\u0441\u043c\u043e\u0442\u0440\u0435\u043b\u0438, \u043a\u0430\u043a \u043d\u0430\u0441\u0442\u0440\u043e\u0438\u0442\u044c \u043a\u0430\u0441\u0442\u043e\u043c\u043d\u044b\u0439 TimestampExtractor, \u0447\u0442\u043e\u0431\u044b \u043a\u043e\u0440\u0440\u0435\u043a\u0442\u043d\u043e \u043e\u0431\u0440\u0430\u0431\u0430\u0442\u044b\u0432\u0430\u0442\u044c \u0441\u043e\u0431\u044b\u0442\u0438\u044f \u0434\u0430\u0436\u0435 \u043f\u0440\u0438 \u043d\u0430\u043b\u0438\u0447\u0438\u0438 \u0437\u0430\u0434\u0435\u0440\u0436\u0435\u043a \u0438\u043b\u0438 \u043d\u0435\u0442\u043e\u0447\u043d\u043e\u0441\u0442\u0435\u0439 \u0432\u043e \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u044b\u0445 \u043c\u0435\u0442\u043a\u0430\u0445. \u0420\u0430\u0437\u043e\u0431\u0440\u0430\u043b\u0438\u0441\u044c, \u043a\u0430\u043a \u0440\u0430\u0437\u043b\u0438\u0447\u043d\u044b\u0435 \u0442\u0438\u043f\u044b \u043e\u043a\u043e\u043d (tumbling, hopping, session \u0438 sliding join windows) \u043f\u043e\u0437\u0432\u043e\u043b\u044f\u044e\u0442 \u044d\u0444\u0444\u0435\u043a\u0442\u0438\u0432\u043d\u043e \u0433\u0440\u0443\u043f\u043f\u0438\u0440\u043e\u0432\u0430\u0442\u044c \u0438 \u0430\u043d\u0430\u043b\u0438\u0437\u0438\u0440\u043e\u0432\u0430\u0442\u044c \u0434\u0430\u043d\u043d\u044b\u0435 \u043f\u043e \u0432\u044b\u0431\u0440\u0430\u043d\u043d\u044b\u043c \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u0430\u043c \u0432\u0440\u0435\u043c\u0435\u043d\u0438, \u0438 \u043a\u0430\u043a \u043e\u043f\u0435\u0440\u0430\u0442\u043e\u0440 suppress \u043f\u043e\u043c\u043e\u0433\u0430\u0435\u0442 \u043d\u0435 \u043f\u0443\u0431\u043b\u0438\u043a\u043e\u0432\u0430\u0442\u044c \u043f\u0440\u043e\u043c\u0435\u0436\u0443\u0442\u043e\u0447\u043d\u044b\u0435 \u0440\u0435\u0437\u0443\u043b\u044c\u0442\u0430\u0442\u044b \u0434\u043e \u0437\u0430\u0432\u0435\u0440\u0448\u0435\u043d\u0438\u044f \u043e\u043a\u043d\u0430. \u041d\u0430\u043a\u043e\u043d\u0435\u0446, \u043c\u044b \u0448\u0430\u0433 \u0437\u0430 \u0448\u0430\u0433\u043e\u043c \u0441\u043e\u0437\u0434\u0430\u043b\u0438 \u043f\u043e\u043b\u043d\u043e\u0446\u0435\u043d\u043d\u043e\u0435 \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0435 \u043d\u0430 Kafka Streams, \u0441\u043f\u043e\u0441\u043e\u0431\u043d\u043e\u0435 \u043f\u0440\u0438\u043d\u0438\u043c\u0430\u0442\u044c, \u043e\u0431\u0440\u0430\u0431\u0430\u0442\u044b\u0432\u0430\u0442\u044c \u0438 \u043e\u0431\u044a\u0435\u0434\u0438\u043d\u044f\u0442\u044c \u043f\u043e\u0442\u043e\u043a\u043e\u0432\u044b\u0435 \u0434\u0430\u043d\u043d\u044b\u0435 \u043e \u0436\u0438\u0437\u043d\u0435\u043d\u043d\u043e \u0432\u0430\u0436\u043d\u044b\u0445 \u043f\u043e\u043a\u0430\u0437\u0430\u0442\u0435\u043b\u044f\u0445 (\u043f\u0443\u043b\u044c\u0441\u0435 \u0438 \u0442\u0435\u043c\u043f\u0435\u0440\u0430\u0442\u0443\u0440\u0435 \u0442\u0435\u043b\u0430), \u0434\u0435\u043c\u043e\u043d\u0441\u0442\u0440\u0438\u0440\u0443\u044f, \u043a\u0430\u043a \u0432\u0441\u0435 \u044d\u0442\u0438 \u0438\u043d\u0441\u0442\u0440\u0443\u043c\u0435\u043d\u0442\u044b \u0440\u0430\u0431\u043e\u0442\u0430\u044e\u0442 \u0432\u043c\u0435\u0441\u0442\u0435 \u0434\u043b\u044f \u0440\u0435\u0448\u0435\u043d\u0438\u044f \u0440\u0435\u0430\u043b\u044c\u043d\u044b\u0445 \u0437\u0430\u0434\u0430\u0447 \u043c\u043e\u043d\u0438\u0442\u043e\u0440\u0438\u043d\u0433\u0430.<\/p>\n<p>\u0414\u0430\u043d\u043d\u044b\u0439 \u043a\u043e\u0434 \u043c\u043e\u0436\u043d\u043e \u0440\u0430\u0441\u0448\u0438\u0440\u044f\u0442\u044c, \u0434\u043e\u0431\u0430\u0432\u043b\u044f\u044f \u043d\u043e\u0432\u044b\u0435 \u043c\u0435\u0442\u0440\u0438\u043a\u0438 (\u0434\u0430\u0432\u043b\u0435\u043d\u0438\u0435, \u0443\u0440\u043e\u0432\u0435\u043d\u044c \u043a\u0438\u0441\u043b\u043e\u0440\u043e\u0434\u0430 \u0432 \u043a\u0440\u043e\u0432\u0438 \u0438 \u0442.\u0434.) \u0438\u043b\u0438 \u0432\u043d\u0435\u0434\u0440\u044f\u044f \u043b\u043e\u0433\u0438\u043a\u0443 \u043e\u043f\u043e\u0432\u0435\u0449\u0435\u043d\u0438\u0439 \u043f\u0440\u0438 \u0432\u044b\u0445\u043e\u0434\u0435 \u043f\u043e\u043a\u0430\u0437\u0430\u0442\u0435\u043b\u0435\u0439 \u0437\u0430 \u0434\u043e\u043f\u0443\u0441\u0442\u0438\u043c\u044b\u0435 \u043f\u0440\u0435\u0434\u0435\u043b\u044b, \u0432 \u043e\u0431\u0449\u0435\u043c \u0435\u0441\u0442\u044c \u043d\u0430\u0434 \u0447\u0435\u043c \u0435\u0449\u0435 \u0444\u0430\u043d\u0442\u0430\u0437\u0438\u0440\u043e\u0432\u0430\u0442\u044c =)<\/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\/870784\/\"> https:\/\/habr.com\/ru\/articles\/870784\/<\/a><\/p>\n","protected":false},"excerpt":{"rendered":"<div><!--[--><!--]--><\/div>\n<div id=\"post-content-body\">\n<div>\n<div class=\"article-formatted-body article-formatted-body article-formatted-body_version-2\">\n<div xmlns=\"http:\/\/www.w3.org\/1999\/xhtml\">\n<p>\u042d\u0442\u043e \u043c\u043e\u044f \u0444\u0438\u043d\u0430\u043b\u044c\u043d\u0430\u044f \u0447\u0430\u0441\u0442\u044c(\u043d\u0443 \u043f\u043e\u043a\u0430 \u0447\u0442\u043e ;)) \u0441\u0435\u0440\u0438\u0438 \u0441\u0442\u0430\u0442\u0435\u0439 \u043f\u0440\u043e Kafka Streams, \u043f\u0440\u043e\u0448\u043b\u044b\u0435 \u0441\u0442\u0430\u0442\u044c\u0438 \u0442\u0443\u0442 [<a href=\"https:\/\/habr.com\/ru\/articles\/850832\/\" rel=\"noopener noreferrer nofollow\">\u043d\u043e\u043b\u044c<\/a>, <a href=\"https:\/\/habr.com\/ru\/articles\/854680\/\" rel=\"noopener noreferrer nofollow\">\u043e\u0434\u0438\u043d<\/a>, <a href=\"https:\/\/habr.com\/ru\/articles\/858668\/\" rel=\"noopener noreferrer nofollow\">\u0434\u0432\u0430<\/a>, <a href=\"https:\/\/habr.com\/ru\/articles\/862976\/\" rel=\"noopener noreferrer nofollow\">\u0442\u0440\u0438<\/a>] \u0422\u0435\u043f\u0435\u0440\u044c \u0434\u0430\u0432\u0430\u0439\u0442\u0435 \u0440\u0430\u0437\u0440\u0430\u0431\u043e\u0442\u0430\u0435\u043c \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0435, \u043a\u043e\u0442\u043e\u0440\u043e\u0435 \u0441\u0447\u0438\u0442\u044b\u0432\u0430\u0435\u0442:<\/p>\n<p>\u0421\u043e\u0431\u044b\u0442\u0438\u044f \u043e \u043f\u0443\u043b\u044c\u0441\u0435 \u0438\u0437 \u0442\u043e\u043f\u0438\u043a\u0430 <strong>pulse-events<\/strong>.<\/p>\n<pre><code class=\"json\">{ \"timestamp\": \"2024-12-26T09:02:00.000Z\" }<\/code><\/pre>\n<p>\u0421\u043e\u0431\u044b\u0442\u0438\u044f \u043e \u0442\u0435\u043c\u043f\u0435\u0440\u0430\u0442\u0443\u0440\u0435 \u0442\u0435\u043b\u0430 \u0438\u0437 \u0442\u043e\u043f\u0438\u043a\u0430 <strong>body-temp-events<\/strong>.<\/p>\n<pre><code class=\"json\">{ \"timestamp\": \"2024-12-26T09:02:00.000Z\", \"temperature\": 37.4, \"unit\": \"\u0421\" }<\/code><\/pre>\n<p>\u041f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0435 \u0431\u0443\u0434\u0435\u0442 \u043e\u0431\u0440\u0430\u0431\u0430\u0442\u044b\u0432\u0430\u0442\u044c \u044d\u0442\u0438 \u0441\u043e\u0431\u044b\u0442\u0438\u044f, \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u044f <strong>Kafka Streams<\/strong>. \u041c\u044b \u0431\u0443\u0434\u0435\u043c \u043e\u0440\u0438\u0435\u043d\u0442\u0438\u0440\u043e\u0432\u0430\u0442\u044c\u0441\u044f \u043d\u0430 <strong>event time<\/strong> (\u0432\u0440\u0435\u043c\u044f, \u0441\u043e\u0434\u0435\u0440\u0436\u0430\u0449\u0435\u0435\u0441\u044f \u0432 \u0441\u0430\u043c\u043e\u043c \u0441\u043e\u0431\u044b\u0442\u0438\u0438), \u0430 \u043d\u0435 \u043d\u0430 \u0432\u0440\u0435\u043c\u044f \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438 \u0438\u043b\u0438 \u0432\u0440\u0435\u043c\u044f \u0437\u0430\u0433\u0440\u0443\u0437\u043a\u0438 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f \u0432 Kafka.<\/p>\n<hr\/>\n<p><strong>\u041f\u043e\u043d\u044f\u0442\u0438\u044f \u0432\u0440\u0435\u043c\u0435\u043d\u0438 \u0432 Kafka Streams<\/strong><\/p>\n<p>\u0412 \u043f\u043e\u0442\u043e\u043a\u043e\u0432\u043e\u0439 \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0435 \u043e\u0431\u044b\u0447\u043d\u043e \u0432\u044b\u0434\u0435\u043b\u044f\u044e\u0442 \u0442\u0440\u0438 \u0432\u0438\u0434\u0430 \u0432\u0440\u0435\u043c\u0435\u043d\u0438:<\/p>\n<ul>\n<li>\n<p><strong>Event Time<\/strong> (\u0432\u0440\u0435\u043c\u044f \u0441\u043e\u0431\u044b\u0442\u0438\u044f). \u041a\u043e\u0433\u0434\u0430 \u0441\u043e\u0431\u044b\u0442\u0438\u0435 \u0440\u0435\u0430\u043b\u044c\u043d\u043e \u043f\u0440\u043e\u0438\u0437\u043e\u0448\u043b\u043e (\u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0435 \u043f\u043e\u043b\u044f timestamp, \u043f\u0440\u043e\u0441\u0442\u0430\u0432\u043b\u0435\u043d\u043d\u043e\u0435 \u0434\u0430\u0442\u0447\u0438\u043a\u043e\u043c \u0438\u043b\u0438 \u0441\u0435\u0440\u0432\u0438\u0441\u043e\u043c).<\/p>\n<\/li>\n<li>\n<p><strong>Ingestion Time<\/strong> (\u0438\u043b\u0438 Loading Time). \u0412\u0440\u0435\u043c\u044f, \u043a\u043e\u0433\u0434\u0430 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0435 \u043f\u043e\u043f\u0430\u0434\u0430\u0435\u0442 \u0432 \u0441\u0438\u0441\u0442\u0435\u043c\u0443 (Kafka), \u0444\u0438\u043a\u0441\u0438\u0440\u0443\u0435\u0442\u0441\u044f \u0431\u0440\u043e\u043a\u0435\u0440\u043e\u043c Kafka.<\/p>\n<\/li>\n<li>\n<p><strong>Processing Time<\/strong> (\u0432\u0440\u0435\u043c\u044f \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438). \u0422\u0435\u043a\u0443\u0449\u0435\u0435 \u0441\u0438\u0441\u0442\u0435\u043c\u043d\u043e\u0435 \u0432\u0440\u0435\u043c\u044f \u043f\u0440\u0438 \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0435 \u0441\u043e\u0431\u044b\u0442\u0438\u044f \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0435\u043c (JVM time).<\/p>\n<\/li>\n<\/ul>\n<p>\u0418\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u044f <strong>event time<\/strong>, \u0447\u0442\u043e\u0431\u044b \u0430\u043d\u0430\u043b\u0438\u0437\u0438\u0440\u043e\u0432\u0430\u0442\u044c \u0441\u043e\u0431\u044b\u0442\u0438\u044f \u0432 \u0442\u043e\u043c \u043f\u043e\u0440\u044f\u0434\u043a\u0435, \u0432 \u043a\u0430\u043a\u043e\u043c \u043e\u043d\u0438 \u0441\u043b\u0443\u0447\u0438\u043b\u0438\u0441\u044c \u0432 \u0440\u0435\u0430\u043b\u044c\u043d\u043e\u043c \u043c\u0438\u0440\u0435. \u0415\u0441\u043b\u0438 \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u044b\u0435 \u043c\u0435\u0442\u043a\u0438 \u043d\u0435\u0432\u0435\u0440\u043d\u044b\u0435 \u0438\u043b\u0438 \u043e\u0442\u0441\u0443\u0442\u0441\u0442\u0432\u0443\u044e\u0442, \u043c\u043e\u0436\u043d\u043e \u043e\u0440\u0438\u0435\u043d\u0442\u0438\u0440\u043e\u0432\u0430\u0442\u044c\u0441\u044f \u043d\u0430 \u0432\u0440\u0435\u043c\u044f \u0437\u0430\u0433\u0440\u0443\u0437\u043a\u0438 \u0438\u043b\u0438 \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438.<\/p>\n<hr\/>\n<p><strong>\u041f\u043e\u0434\u0433\u043e\u0442\u043e\u0432\u043a\u0430 \u043e\u043a\u0440\u0443\u0436\u0435\u043d\u0438\u044f \u0438 \u0441\u0442\u0440\u0443\u043a\u0442\u0443\u0440\u0430 \u043f\u0440\u043e\u0435\u043a\u0442\u0430<\/strong><\/p>\n<p><strong>\u041d\u0435\u043e\u0431\u0445\u043e\u0434\u0438\u043c\u044b\u0435 \u0438\u043d\u0441\u0442\u0440\u0443\u043c\u0435\u043d\u0442\u044b<\/strong><\/p>\n<ul>\n<li>\n<p><strong>Apache Kafka<\/strong> (\u0437\u0430\u043f\u0443\u0449\u0435\u043d\u043d\u0430\u044f \u043b\u043e\u043a\u0430\u043b\u044c\u043d\u043e \u0438\u043b\u0438 \u0432 Docker)<\/p>\n<\/li>\n<li>\n<p><strong>Gradle<\/strong> (\u043c\u044b \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u043c Kotlin DSL \u2014 build.gradle.kts)<\/p>\n<\/li>\n<li>\n<p><strong>Java 21<\/strong><\/p>\n<\/li>\n<\/ul>\n<p>\u0423\u0431\u0435\u0434\u0438\u0442\u0435\u0441\u044c, \u0447\u0442\u043e Kafka \u0437\u0430\u043f\u0443\u0449\u0435\u043d\u0430 \u043d\u0430 localhost:9092.<\/p>\n<p><strong>\u0421\u0442\u0440\u0443\u043a\u0442\u0443\u0440\u0430 \u043f\u0440\u043e\u0435\u043a\u0442\u0430<\/strong><\/p>\n<p><em>kafka-streams-medical-app<\/em><\/p>\n<p>\u251c\u2500 src<\/p>\n<p>\u2502\u00a0 \u251c\u2500 main<\/p>\n<p>\u2502\u00a0 \u2502\u00a0 \u251c\u2500 java<\/p>\n<p>\u2502\u00a0 \u2502\u00a0 \u2502\u00a0 \u2514\u2500 com.example.kafka<\/p>\n<p>\u2502\u00a0 \u2502\u00a0 \u2502 \u00a0 \u00a0 \u251c\u2500 MedicalMonitorApp.java<\/p>\n<p>\u2502\u00a0 \u2502\u00a0 \u2502 \u00a0 \u00a0 \u251c\u2500 PulseEvent.java<\/p>\n<p>\u2502\u00a0 \u2502\u00a0 \u2502 \u00a0 \u00a0 \u251c\u2500 BodyTempEvent.java<\/p>\n<p>\u2502\u00a0 \u2502\u00a0 \u2502 \u00a0 \u00a0 \u251c\u2500 MyTimestampExtractor.java<\/p>\n<p>\u2502\u00a0 \u2502\u00a0 \u2502 \u00a0 \u00a0 \u2514\u2500 &#8230;<\/p>\n<p>\u2502\u00a0 \u2502\u00a0 \u2514\u2500 resources<\/p>\n<p>\u2502\u00a0 \u2514\u2500 test<\/p>\n<p>\u2502 \u00a0 \u00a0 \u2514\u2500 java<\/p>\n<p>\u251c\u2500 build.gradle.kts<\/p>\n<p>\u2514\u2500 settings.gradle.kts<\/p>\n<p><strong>\u041f\u0440\u0438\u043c\u0435\u0440 build.gradle.kts<\/strong><\/p>\n<p>\u0423\u043f\u0440\u043e\u0449\u0451\u043d\u043d\u0430\u044f \u0432\u0435\u0440\u0441\u0438\u044f \u0441\u0431\u043e\u0440\u043e\u0447\u043d\u043e\u0433\u043e \u0444\u0430\u0439\u043b\u0430 (Kotlin DSL):<\/p>\n<pre><code class=\"java\">plugins { kotlin(\"jvm\") application  }  repositories { mavenCentral() }  dependencies { implementation(\"org.apache.kafka:kafka-streams:3.8.0\") implementation(\"com.fasterxml.jackson.core:jackson-databind:2.15.0\") implementation(\"com.fasterxml.jackson.module:jackson-module-kotlin:2.15.0\") implementation(\"org.slf4j:slf4j-api:2.0.0\")  runtimeOnly(\"org.slf4j:slf4j-simple:2.0.0\")  testImplementation(kotlin(\"test\")) testImplementation(\"org.apache.kafka:kafka-streams-test-utils:3.8.0\") }  application { mainClass.set(\"com.example.kafka.MedicalMonitorApp\") }<\/code><\/pre>\n<p>\u0417\u0430\u043f\u0443\u0441\u043a \u0441\u0431\u043e\u0440\u043a\u0438: <code>.\/gradlew clean build<\/code><\/p>\n<hr\/>\n<p><strong>\u0421\u043e\u0437\u0434\u0430\u043d\u0438\u0435 \u0438 \u043d\u0430\u0441\u0442\u0440\u043e\u0439\u043a\u0430 \u0442\u043e\u043f\u0438\u043a\u043e\u0432<\/strong><\/p>\n<p>\u0414\u043b\u044f \u0440\u0430\u0431\u043e\u0442\u044b \u043d\u0443\u0436\u043d\u043e \u0441\u043e\u0437\u0434\u0430\u0442\u044c \u0434\u0432\u0430 \u0442\u043e\u043f\u0438\u043a\u0430:<\/p>\n<ul>\n<li>\n<p><strong>pulse-events<\/strong><\/p>\n<\/li>\n<li>\n<p><strong>body-temp-events<\/strong><\/p>\n<\/li>\n<\/ul>\n<p>\u041f\u0440\u0438\u043c\u0435\u0440 \u043a\u043e\u043c\u0430\u043d\u0434 (\u043a\u043e\u043d\u0441\u043e\u043b\u044c Kafka):<\/p>\n<pre><code class=\"bash\">bin\/kafka-topics.sh --bootstrap-server localhost:9092 \\\\ --create --topic pulse-events \\\\ --partitions 1 --replication-factor 1  bin\/kafka-topics.sh --bootstrap-server localhost:9092 \\\\ --create --topic body-temp-events \\\\ --partitions 1 --replication-factor 1<\/code><\/pre>\n<p><strong>\u0421\u043e\u0437\u0434\u0430\u043d\u0438\u0435 \u043a\u043b\u0430\u0441\u0441\u043e\u0432-\u0441\u0443\u0449\u043d\u043e\u0441\u0442\u0435\u0439 (DTO)<\/strong><\/p>\n<p>\u041f\u0440\u0438\u043c\u0435\u0440 \u043a\u043b\u0430\u0441\u0441\u0430 <strong>PulseEvent.java<\/strong>:<\/p>\n<pre><code class=\"java\">public record PulseEvent(String timestamp) { }<\/code><\/pre>\n<p>\u041f\u0440\u0438\u043c\u0435\u0440 \u043a\u043b\u0430\u0441\u0441\u0430 <strong>BodyTempEvent.java<\/strong>:<\/p>\n<pre><code class=\"java\">public record BodyTempEvent(String timestamp, double temperature, String unit) { }<\/code><\/pre>\n<hr\/>\n<p><strong>\u041a\u0430\u043a \u0440\u0430\u0437\u043b\u0438\u0447\u0430\u0442\u044c \u0432\u0440\u0435\u043c\u044f \u0441\u043e\u0431\u044b\u0442\u0438\u044f, \u0432\u0440\u0435\u043c\u044f \u0437\u0430\u0433\u0440\u0443\u0437\u043a\u0438 \u0438 \u0432\u0440\u0435\u043c\u044f \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438<\/strong><\/p>\n<p>\u041c\u044b \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u043c <strong>event time<\/strong> (\u0438\u0437 \u043f\u043e\u043b\u044f timestamp \u0432\u043e \u0432\u0445\u043e\u0434\u044f\u0449\u0435\u043c \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0438), \u0447\u0442\u043e\u0431\u044b \u0430\u043d\u0430\u043b\u0438\u0437\u0438\u0440\u043e\u0432\u0430\u0442\u044c \u0441\u043e\u0431\u044b\u0442\u0438\u044f \u0432 \u0440\u0435\u0430\u043b\u044c\u043d\u043e\u043c \u043f\u043e\u0440\u044f\u0434\u043a\u0435 \u0438\u0445 \u0432\u043e\u0437\u043d\u0438\u043a\u043d\u043e\u0432\u0435\u043d\u0438\u044f. Kafka \u043f\u043e \u0443\u043c\u043e\u043b\u0447\u0430\u043d\u0438\u044e \u043c\u043e\u0436\u0435\u0442 \u0431\u0440\u0430\u0442\u044c \u0432\u0440\u0435\u043c\u044f \u0437\u0430\u0433\u0440\u0443\u0437\u043a\u0438 (ingestion), \u043d\u043e \u0434\u043b\u044f \u043c\u0435\u0434\u0438\u0446\u0438\u043d\u0441\u043a\u0438\u0445 \u043f\u043e\u043a\u0430\u0437\u0430\u0442\u0435\u043b\u0435\u0439 \u0432\u0430\u0436\u043d\u043e \u0443\u0447\u0438\u0442\u044b\u0432\u0430\u0442\u044c \u0444\u0430\u043a\u0442\u0438\u0447\u0435\u0441\u043a\u043e\u0435 \u0432\u0440\u0435\u043c\u044f \u0438\u0437\u043c\u0435\u0440\u0435\u043d\u0438\u044f.<\/p>\n<p>\u0415\u0441\u043b\u0438 \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u044b\u0435 \u043c\u0435\u0442\u043a\u0438 \u043d\u0435\u043a\u043e\u0440\u0440\u0435\u043a\u0442\u043d\u044b (\u0438\u043b\u0438 \u043e\u0442\u0441\u0443\u0442\u0441\u0442\u0432\u0443\u044e\u0442), \u043c\u043e\u0436\u043d\u043e \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c fallback \u043d\u0430 ingestion time \u0438\u043b\u0438 \u043d\u0430 processing time (\u0441\u0438\u0441\u0442\u0435\u043c\u043d\u043e\u0435 \u0432\u0440\u0435\u043c\u044f \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438).<\/p>\n<hr\/>\n<p><strong>\u0421\u043e\u0437\u0434\u0430\u043d\u0438\u0435 \u043a\u0430\u0441\u0442\u043e\u043c\u043d\u043e\u0433\u043e TimestampExtractor<\/strong><\/p>\n<p>\u0427\u0442\u043e\u0431\u044b Kafka Streams \u0437\u043d\u0430\u043b, \u0447\u0442\u043e \u043c\u044b \u043e\u0440\u0438\u0435\u043d\u0442\u0438\u0440\u0443\u0435\u043c\u0441\u044f \u043d\u0430 <strong>event time<\/strong>, \u043d\u0443\u0436\u043d\u043e \u043d\u0430\u0441\u0442\u0440\u043e\u0438\u0442\u044c TimestampExtractor, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u0431\u0443\u0434\u0435\u0442 \u0437\u0430\u0431\u0438\u0440\u0430\u0442\u044c \u043f\u043e\u043b\u0435 timestamp \u0438\u0437 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f:<\/p>\n<pre><code class=\"java\">import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.streams.processor.TimestampExtractor; import com.fasterxml.jackson.databind.ObjectMapper;  public class MyTimestampExtractor implements TimestampExtractor {      private final ObjectMapper objectMapper = new ObjectMapper();      @Override     public long extract(final ConsumerRecord&lt;Object, Object&gt; record, final long partitionTime) {         if (record.value() == null) {             return partitionTime;         }          try {             String jsonString = record.value().toString();             String timestampStr = objectMapper.readTree(jsonString).get(\"timestamp\").asText();             return javax.xml.bind.DatatypeConverter.parseDateTime(timestampStr).getTimeInMillis();         } catch (Exception e) {             return partitionTime;         }     } }<\/code><\/pre>\n<ul>\n<li>\n<p>\u0418\u0437\u0432\u043b\u0435\u043a\u0430\u0435\u043c JSON-\u043f\u043e\u043b\u0435 &#171;timestamp&#187; \u0438\u0437 \u0441\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u044f.<\/p>\n<\/li>\n<li>\n<p>\u041f\u0440\u0435\u043e\u0431\u0440\u0430\u0437\u0443\u0435\u043c \u0435\u0433\u043e \u0438\u0437 ISO-8601 \u0432 \u043c\u0438\u043b\u043b\u0438\u0441\u0435\u043a\u0443\u043d\u0434\u044b (long).<\/p>\n<\/li>\n<li>\n<p>\u0415\u0441\u043b\u0438 \u043f\u0430\u0440\u0441\u0438\u043d\u0433 \u043d\u0435 \u0443\u0434\u0430\u043b\u0441\u044f, \u0432\u043e\u0437\u0432\u0440\u0430\u0449\u0430\u0435\u043c partitionTime.<\/p>\n<\/li>\n<\/ul>\n<hr\/>\n<p><strong>\u041a\u0430\u043a \u0432\u0440\u0435\u043c\u044f \u0443\u043f\u0440\u0430\u0432\u043b\u044f\u0435\u0442 \u043f\u043e\u0442\u043e\u043a\u043e\u043c \u0434\u0430\u043d\u043d\u044b\u0445 \u0432 Kafka Streams<\/strong><\/p>\n<p>\u0411\u043b\u0430\u0433\u043e\u0434\u0430\u0440\u044f \u043a\u0430\u0441\u0442\u043e\u043c\u043d\u043e\u043c\u0443 TimestampExtractor, Kafka Streams \u0432\u044b\u0441\u0442\u0440\u0430\u0438\u0432\u0430\u0435\u0442 \u0441\u043e\u0431\u044b\u0442\u0438\u044f \u043f\u043e \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u043e\u0439 \u0448\u043a\u0430\u043b\u0435 \u043d\u0430 \u043e\u0441\u043d\u043e\u0432\u0435 <strong>event time<\/strong>. \u0415\u0441\u043b\u0438 \u043f\u0440\u0438\u0445\u043e\u0434\u044f\u0442 \u201c\u0437\u0430\u043f\u0430\u0437\u0434\u044b\u0432\u0430\u044e\u0449\u0438\u0435\u201d \u0441\u043e\u0431\u044b\u0442\u0438\u044f \u0441\u043e \u0441\u0442\u0430\u0440\u044b\u043c\u0438 \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u044b\u043c\u0438 \u043c\u0435\u0442\u043a\u0430\u043c\u0438, \u043e\u043d\u0438 \u043c\u043e\u0433\u0443\u0442 \u043f\u043e\u043f\u0430\u0441\u0442\u044c \u0432 \u0443\u0436\u0435 \u0437\u0430\u043a\u0440\u044b\u0442\u044b\u0435 \u043e\u043a\u043d\u0430 (\u044d\u0442\u043e \u0440\u0435\u0433\u0443\u043b\u0438\u0440\u0443\u0435\u0442\u0441\u044f \u043d\u0430\u0441\u0442\u0440\u043e\u0439\u043a\u0430\u043c\u0438 grace period).<\/p>\n<p>\u042d\u0442\u043e \u043e\u0441\u043e\u0431\u0435\u043d\u043d\u043e \u0432\u0430\u0436\u043d\u043e, \u043a\u043e\u0433\u0434\u0430 \u0440\u0435\u0430\u043b\u044c\u043d\u043e\u0435 \u043f\u043e\u043b\u0443\u0447\u0435\u043d\u0438\u0435 \u0441\u043e\u0431\u044b\u0442\u0438\u044f \u0437\u0430\u0434\u0435\u0440\u0436\u0438\u0432\u0430\u0435\u0442\u0441\u044f, \u043d\u043e \u043c\u044b \u0445\u043e\u0442\u0438\u043c \u0443\u0447\u0438\u0442\u044b\u0432\u0430\u0442\u044c \u0435\u0433\u043e \u0438\u043c\u0435\u043d\u043d\u043e \u0432 \u0442\u043e\u043c \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u043e\u043c \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u0435, \u043a \u043a\u043e\u0442\u043e\u0440\u043e\u043c\u0443 \u043e\u043d\u043e \u043f\u0440\u0438\u043d\u0430\u0434\u043b\u0435\u0436\u0438\u0442 \u0444\u0430\u043a\u0442\u0438\u0447\u0435\u0441\u043a\u0438.<\/p>\n<hr\/>\n<p><strong>\u0422\u0438\u043f\u044b \u043e\u043a\u043e\u043d \u0432 Kafka Streams<\/strong><\/p>\n<p><strong>Tumbling Windows<\/strong><\/p>\n<ul>\n<li>\n<p>\u041e\u043a\u043d\u0430 \u0444\u0438\u043a\u0441\u0438\u0440\u043e\u0432\u0430\u043d\u043d\u043e\u0433\u043e \u0440\u0430\u0437\u043c\u0435\u0440\u0430, \u043a\u043e\u0442\u043e\u0440\u044b\u0435 <strong>\u043d\u0435 \u043f\u0435\u0440\u0435\u043a\u0440\u044b\u0432\u0430\u044e\u0442\u0441\u044f<\/strong>.<\/p>\n<\/li>\n<li>\n<p>\u041d\u0430\u043f\u0440\u0438\u043c\u0435\u0440, \u043e\u043a\u043d\u043e \u0432 1 \u043c\u0438\u043d\u0443\u0442\u0443: [09:00 \u2013 09:01], [09:01 \u2013 09:02] \u0438 \u0442.\u0434.<\/p>\n<\/li>\n<\/ul>\n<p><strong>Hopping Windows<\/strong><\/p>\n<ul>\n<li>\n<p>\u041e\u043a\u043d\u0430 \u0444\u0438\u043a\u0441\u0438\u0440\u043e\u0432\u0430\u043d\u043d\u043e\u0433\u043e \u0440\u0430\u0437\u043c\u0435\u0440\u0430, <strong>\u043d\u043e \u043f\u0435\u0440\u0435\u043a\u0440\u044b\u0432\u0430\u044e\u0449\u0438\u0435\u0441\u044f<\/strong>.<\/p>\n<\/li>\n<li>\n<p>\u041f\u0440\u0438\u043c\u0435\u0440: \u0440\u0430\u0437\u043c\u0435\u0440 \u043e\u043a\u043d\u0430 1 \u043c\u0438\u043d\u0443\u0442\u0430, \u0441\u0434\u0432\u0438\u0433 30 \u0441\u0435\u043a\u0443\u043d\u0434, \u043f\u043e\u043b\u0443\u0447\u0430\u0435\u043c [09:00 \u2013 09:01], [09:00:30 \u2013 09:01:30] \u0438 \u0442.\u0434.<\/p>\n<\/li>\n<\/ul>\n<p><strong>Session Windows<\/strong><\/p>\n<ul>\n<li>\n<p>\u041e\u043a\u043d\u0430, \u043a\u043e\u0442\u043e\u0440\u044b\u0435 <strong>\u043d\u0435 \u0438\u043c\u0435\u044e\u0442 \u0444\u0438\u043a\u0441\u0438\u0440\u043e\u0432\u0430\u043d\u043d\u043e\u0439 \u0434\u043b\u0438\u0442\u0435\u043b\u044c\u043d\u043e\u0441\u0442\u0438<\/strong>; \u043e\u043d\u0438 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u044f\u044e\u0442\u0441\u044f \u201c\u043f\u0430\u0443\u0437\u0430\u043c\u0438\u201d \u043c\u0435\u0436\u0434\u0443 \u0441\u043e\u0431\u044b\u0442\u0438\u044f\u043c\u0438.<\/p>\n<\/li>\n<li>\n<p>\u0415\u0441\u043b\u0438 \u043c\u0435\u0436\u0434\u0443 \u0441\u043e\u0431\u044b\u0442\u0438\u044f\u043c\u0438 \u043d\u0435\u0442 \u0431\u043e\u043b\u044c\u0448\u0438\u0445 \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u043e\u0432, \u043e\u043a\u043d\u043e \u0440\u0430\u0441\u0448\u0438\u0440\u044f\u0435\u0442\u0441\u044f; \u043f\u0440\u0438 \u0434\u043b\u0438\u0442\u0435\u043b\u044c\u043d\u043e\u043c \u0437\u0430\u0442\u0438\u0448\u044c\u0435 \u043e\u043a\u043d\u043e \u201c\u0437\u0430\u043a\u0440\u044b\u0432\u0430\u0435\u0442\u0441\u044f\u201d.<\/p>\n<\/li>\n<\/ul>\n<p><strong>Sliding Join Windows<\/strong><\/p>\n<ul>\n<li>\n<p>\u041e\u043a\u043d\u0430, \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u043c\u044b\u0435 \u0432 \u043e\u043f\u0435\u0440\u0430\u0446\u0438\u044f\u0445 join \u043f\u043e\u0442\u043e\u043a\u043e\u0432, \u0433\u0434\u0435 \u043d\u0430\u043c \u0432\u0430\u0436\u043d\u043e \u201c\u0437\u0430\u0445\u0432\u0430\u0442\u0438\u0442\u044c\u201d \u0441\u043e\u0431\u044b\u0442\u0438\u044f, \u043a\u043e\u0442\u043e\u0440\u044b\u0435 \u043f\u0440\u043e\u0438\u0437\u043e\u0448\u043b\u0438 \u0432 \u043f\u0440\u043e\u043c\u0435\u0436\u0443\u0442\u043a\u0435 [t, t+N].<\/p>\n<\/li>\n<li>\n<p>\u041f\u043e\u0445\u043e\u0436\u0435 \u043d\u0430 hopping windows, \u043d\u043e \u043f\u0440\u0438\u043c\u0435\u043d\u044f\u0435\u0442\u0441\u044f \u0434\u043b\u044f Stream-Stream join.<\/p>\n<\/li>\n<\/ul>\n<p><strong>\u0412\u044b\u043f\u043e\u043b\u043d\u0435\u043d\u0438\u0435 \u0438 \u043f\u0440\u043e\u0438\u0437\u0432\u043e\u0434\u0441\u0442\u0432\u043e \u043e\u043a\u043e\u043d\u043d\u044b\u0445 \u0441\u043e\u0435\u0434\u0438\u043d\u0435\u043d\u0438\u0439 (Stream-Stream join)<\/strong><\/p>\n<p>\u0427\u0442\u043e\u0431\u044b \u043e\u0431\u044a\u0435\u0434\u0438\u043d\u0438\u0442\u044c \u0434\u0432\u0430 \u043f\u043e\u0442\u043e\u043a\u0430 \u043d\u0430 \u043e\u0441\u043d\u043e\u0432\u0435 \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u044b\u0445 \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u043e\u0432, \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u0442\u0441\u044f \u043e\u043a\u043e\u043d\u043d\u043e\u0435 \u0441\u043e\u0435\u0434\u0438\u043d\u0435\u043d\u0438\u0435 (Stream-Stream Join). \u041d\u0430\u043f\u0440\u0438\u043c\u0435\u0440, \u0445\u043e\u0442\u0438\u043c \u201c\u0441\u043a\u043b\u0435\u0438\u0442\u044c\u201d \u0434\u0430\u043d\u043d\u044b\u0435 \u043e \u043f\u0443\u043b\u044c\u0441\u0435 \u0438 \u0442\u0435\u043c\u043f\u0435\u0440\u0430\u0442\u0443\u0440\u0435, \u043f\u043e\u0441\u0442\u0443\u043f\u0438\u0432\u0448\u0438\u0435 \u0432 \u043f\u0440\u0435\u0434\u0435\u043b\u0430\u0445 \u043e\u0434\u043d\u043e\u0439 \u043c\u0438\u043d\u0443\u0442\u044b:<\/p>\n<pre><code class=\"java\">KStream&lt;String, PulseEvent&gt; pulseStream = builder     .stream(\"pulse-events\", Consumed.with(Serdes.String(), pulseSerde)         .withTimestampExtractor(new MyTimestampExtractor()));  KStream&lt;String, BodyTempEvent&gt; tempStream = builder     .stream(\"body-temp-events\", Consumed.with(Serdes.String(), bodyTempSerde)         .withTimestampExtractor(new MyTimestampExtractor()));  KStream&lt;String, CombinedMeasurement&gt; joinedStream = pulseStream.join(     tempStream,     (pulseValue, tempValue) -&gt; {         CombinedMeasurement cm = new CombinedMeasurement();         cm.setTimestampPulse(pulseValue.getTimestamp());         cm.setTimestampTemp(tempValue.getTimestamp());         cm.setTemperature(tempValue.getTemperature());         return cm;     },     JoinWindows.ofTimeDifferenceWithNoGrace(Duration.ofMinutes(1)),     StreamJoined.with(Serdes.String(), pulseSerde, bodyTempSerde) );<\/code><\/pre>\n<p>\u0422\u0430\u043a\u0438\u043c \u043e\u0431\u0440\u0430\u0437\u043e\u043c, \u0435\u0441\u043b\u0438 \u0441\u043e\u0431\u044b\u0442\u0438\u044f \u043f\u043e \u043a\u043b\u044e\u0447\u0443 \u0441\u043e\u0432\u043f\u0430\u0434\u0430\u044e\u0442 \u0438 \u043f\u043e\u043f\u0430\u0434\u0430\u044e\u0442 \u0432 \u043e\u0431\u0449\u0438\u0439 \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u043e\u0439 \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b, \u043e\u043d\u0438 \u0431\u0443\u0434\u0443\u0442 \u043e\u0431\u044a\u0435\u0434\u0438\u043d\u0435\u043d\u044b \u0432 \u0435\u0434\u0438\u043d\u044b\u0439 \u0440\u0435\u0437\u0443\u043b\u044c\u0442\u0430\u0442.<\/p>\n<hr\/>\n<p><strong>\u0417\u0430\u0447\u0435\u043c \u043d\u0443\u0436\u0435\u043d \u043e\u043f\u0435\u0440\u0430\u0442\u043e\u0440 suppress \u0438 \u043a\u0430\u043a \u0435\u0433\u043e \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c<\/strong><\/p>\n<p>\u0412 KTable \u0438 \u0430\u0433\u0440\u0435\u0433\u0438\u0440\u0443\u044e\u0449\u0438\u0445 \u043e\u043f\u0435\u0440\u0430\u0446\u0438\u044f\u0445 Kafka Streams \u0432\u043e\u0437\u043d\u0438\u043a\u0430\u044e\u0442 \u043f\u0440\u043e\u043c\u0435\u0436\u0443\u0442\u043e\u0447\u043d\u044b\u0435 \u0440\u0435\u0437\u0443\u043b\u044c\u0442\u0430\u0442\u044b (updates). \u041e\u043f\u0435\u0440\u0430\u0442\u043e\u0440 suppress \u043f\u043e\u0437\u0432\u043e\u043b\u044f\u0435\u0442 \u201c\u0437\u0430\u0434\u0435\u0440\u0436\u0438\u0432\u0430\u0442\u044c\u201d (\u043f\u043e\u0434\u0430\u0432\u043b\u044f\u0442\u044c) \u043f\u0443\u0431\u043b\u0438\u043a\u0430\u0446\u0438\u044e \u043f\u0440\u043e\u043c\u0435\u0436\u0443\u0442\u043e\u0447\u043d\u044b\u0445 \u0440\u0435\u0437\u0443\u043b\u044c\u0442\u0430\u0442\u043e\u0432 \u0438 \u0432\u044b\u0434\u0430\u0432\u0430\u0442\u044c <strong>\u0442\u043e\u043b\u044c\u043a\u043e \u043a\u043e\u043d\u0435\u0447\u043d\u044b\u0439<\/strong> (\u043f\u043e\u0441\u043b\u0435 \u0437\u0430\u043a\u0440\u044b\u0442\u0438\u044f \u043e\u043a\u043d\u0430).<\/p>\n<pre><code class=\"java\">KTable&lt;Windowed&lt;String&gt;, Long&gt; aggregatedTable = pulseStream     .groupByKey()     .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(1)))     .count();  KTable&lt;Windowed&lt;String&gt;, Long&gt; suppressedTable = aggregatedTable     .suppress(Suppressed.untilWindowCloses(Suppressed.BufferConfig.unbounded())); <\/code><\/pre>\n<ul>\n<li>\n<p>untilWindowCloses \u043e\u0437\u043d\u0430\u0447\u0430\u0435\u0442, \u0447\u0442\u043e \u043c\u044b \u043f\u0443\u0431\u043b\u0438\u043a\u0443\u0435\u043c \u0440\u0435\u0437\u0443\u043b\u044c\u0442\u0430\u0442 \u0442\u043e\u043b\u044c\u043a\u043e \u043e\u0434\u0438\u043d \u0440\u0430\u0437 \u2014 \u043a\u043e\u0433\u0434\u0430 \u043e\u043a\u043d\u043e \u043e\u043a\u043e\u043d\u0447\u0430\u0442\u0435\u043b\u044c\u043d\u043e \u0437\u0430\u043a\u0440\u044b\u0442\u043e.<\/p>\n<\/li>\n<li>\n<p>\u042d\u0442\u043e \u0438\u0437\u0431\u0430\u0432\u043b\u044f\u0435\u0442 \u043e\u0442 \u201c\u0448\u0442\u043e\u0440\u043c\u0430\u201d \u043f\u043e\u0441\u0442\u043e\u044f\u043d\u043d\u044b\u0445 \u043e\u0431\u043d\u043e\u0432\u043b\u0435\u043d\u0438\u0439.<\/p>\n<\/li>\n<\/ul>\n<hr\/>\n<p><strong>\u0427\u0442\u043e \u0432\u044b\u0448\u043b\u043e \u0432 \u0438\u0442\u043e\u0433\u0435<\/strong><\/p>\n<p>\u041d\u0438\u0436\u0435 \u0443\u043f\u0440\u043e\u0449\u0451\u043d\u043d\u044b\u0439 \u043a\u043e\u0434 <strong>MedicalMonitorApp.java<\/strong>, \u0433\u0434\u0435:<\/p>\n<ul>\n<li>\n<p>\u041d\u0430\u0441\u0442\u0440\u0430\u0438\u0432\u0430\u0435\u043c \u043a\u043e\u043d\u0444\u0438\u0433\u0443\u0440\u0430\u0446\u0438\u044e Kafka Streams.<\/p>\n<\/li>\n<li>\n<p>\u0421\u043e\u0437\u0434\u0430\u0451\u043c \u043f\u043e\u0442\u043e\u043a\u0438 \u0434\u043b\u044f \u043f\u0443\u043b\u044c\u0441\u0430 \u0438 \u0442\u0435\u043c\u043f\u0435\u0440\u0430\u0442\u0443\u0440\u044b.<\/p>\n<\/li>\n<li>\n<p>\u0412\u044b\u043f\u043e\u043b\u043d\u044f\u0435\u043c join \u043f\u043e \u043e\u043a\u043d\u0443 \u0432 1 \u043c\u0438\u043d\u0443\u0442\u0443.<\/p>\n<\/li>\n<li>\n<p>\u041e\u0442\u043f\u0440\u0430\u0432\u043b\u044f\u0435\u043c \u0440\u0435\u0437\u0443\u043b\u044c\u0442\u0430\u0442 \u0432 \u043d\u043e\u0432\u044b\u0439 \u0442\u043e\u043f\u0438\u043a.<\/p>\n<\/li>\n<\/ul>\n<pre><code class=\"java\">import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.*; import org.apache.kafka.streams.kstream.*; import java.time.Duration; import java.util.Properties;  public class MedicalMonitorApp {      public static void main(String[] args) {         Properties props = new Properties();         props.put(StreamsConfig.APPLICATION_ID_CONFIG, \"medical-monitor-app\");         props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, \"localhost:9092\");         props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.StringSerde.class);         props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.StringSerde.class);         props.put(StreamsConfig.DEFAULT_TIMESTAMP_EXTRACTOR_CLASS_CONFIG, MyTimestampExtractor.class);          StreamsBuilder builder = new StreamsBuilder();          KStream&lt;String, PulseEvent&gt; pulseStream = builder.stream(             \"pulse-events\",             Consumed.with(                 Serdes.String(),                 Serdes.serdeFrom(new JsonSerializer&lt;&gt;(), new JsonDeserializer&lt;&gt;(PulseEvent.class))             )         );          KStream&lt;String, BodyTempEvent&gt; tempStream = builder.stream(             \"body-temp-events\",             Consumed.with(                 Serdes.String(),                 Serdes.serdeFrom(new JsonSerializer&lt;&gt;(), new JsonDeserializer&lt;&gt;(BodyTempEvent.class))             )         );          KStream&lt;String, String&gt; joinedStream = pulseStream.join(             tempStream,             (pulseVal, tempVal) -&gt; {                 return \"Pulse timestamp: \" + pulseVal.getTimestamp() +                        \", Temp timestamp: \" + tempVal.getTimestamp() +                        \", Temp: \" + tempVal.getTemperature();             },             JoinWindows.ofTimeDifferenceWithNoGrace(Duration.ofMinutes(1)),        <\/code><\/pre>\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-443829","post","type-post","status-publish","format-standard","hentry"],"_links":{"self":[{"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=\/wp\/v2\/posts\/443829","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=443829"}],"version-history":[{"count":0,"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=\/wp\/v2\/posts\/443829\/revisions"}],"wp:attachment":[{"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=%2Fwp%2Fv2%2Fmedia&parent=443829"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=%2Fwp%2Fv2%2Fcategories&post=443829"},{"taxonomy":"post_tag","embeddable":true,"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=%2Fwp%2Fv2%2Ftags&post=443829"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}