{"id":35786,"date":"2019-10-31T22:06:19","date_gmt":"2019-10-31T19:06:19","guid":{"rendered":"https:\/\/prohoster.info\/blog\/kniga-kafka-streams-v-dejstvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni\/"},"modified":"2019-10-31T22:06:19","modified_gmt":"2019-10-31T19:06:19","slug":"kniga-kafka-streams-v-dejstvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni","status":"publish","type":"post","link":"https:\/\/prohoster.info\/es\/blog\/administrirovanie\/kniga-kafka-streams-v-dejstvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni","title":{"rendered":"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb","gt_translate_keys":[{"key":"rendered","format":"text"}]},"content":{"rendered":"<p><noindex><a rel=\"nofollow\" href=\"https:\/\/habr.com\/ru\/company\/piter\/blog\/457756\/\"><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/6f8bd2b31b87b0c760c1148515893c43.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/a><\/noindex> \u00a1Hola, habitantes de Habr! Este libro es adecuado para cualquier desarrollador que quiera entender el procesamiento de flujos. Comprender la programaci\u00f3n distribuida ayudar\u00e1 a estudiar mejor Kafka y Kafka Streams. Ser\u00eda bueno conocer el propio framework Kafka, pero no es obligatorio: te contar\u00e9 todo lo que necesitas. Los desarrolladores experimentados en Kafka, as\u00ed como los principiantes, dominar\u00e1n la creaci\u00f3n de aplicaciones interesantes para el procesamiento de flujos utilizando la biblioteca Kafka Streams, gracias a este libro. Los desarrolladores de Java de nivel medio y alto, ya familiarizados con conceptos como la serializaci\u00f3n, aprender\u00e1n a aplicar sus habilidades para crear aplicaciones de Kafka Streams. El c\u00f3digo fuente del libro est\u00e1 escrito en Java 8 y utiliza de manera significativa la sintaxis de expresiones lambda de Java 8, por lo que la capacidad de trabajar con funciones lambda (incluso en otro lenguaje de programaci\u00f3n) ser\u00e1 \u00fatil.<br \/>\n<noindex><a rel=\"nofollow\" name=\"habracut\"><\/a><\/noindex><\/p>\n<h3>Fragmento. 5.3. Agregaci\u00f3n y operaciones de ventana<\/h3>\n<p>\nEn esta secci\u00f3n, pasaremos a estudiar las partes m\u00e1s prometedoras de Kafka Streams. Hasta ahora hemos examinado los siguientes aspectos de Kafka Streams:<\/p>\n<ul>\n<li>creaci\u00f3n de topolog\u00edas de procesamiento;<\/li>\n<li>uso de estado en aplicaciones de flujo;<\/li>\n<li>realizaci\u00f3n de uniones de flujos de datos;<\/li>\n<li>diferencias entre flujos de eventos (KStream) y flujos de actualizaciones (KTable).<\/li>\n<\/ul>\n<p>\nEn los siguientes ejemplos, reuniremos todos estos elementos. Adem\u00e1s, te familiarizar\u00e1s con las operaciones de ventana, otra gran capacidad de las aplicaciones de flujo. Nuestro primer ejemplo ser\u00e1 la agregaci\u00f3n simple.<\/p>\n<h3>5.3.1. Agregaci\u00f3n del volumen de ventas de acciones por sectores industriales<\/h3>\n<p>\nLa agregaci\u00f3n y agrupaci\u00f3n son herramientas vitales al trabajar con datos de flujo. Examinar registros individuales a medida que llegan a menudo resulta insuficiente. Para extraer informaci\u00f3n adicional de los datos, es necesario agrupamiento y combinaci\u00f3n.<\/p>\n<p>En este ejemplo, te pondr\u00e1s en la piel de un trader intrad\u00eda que necesita rastrear los vol\u00famenes de ventas de acciones de empresas en varios sectores industriales. En particular, te interesan las cinco empresas con los mayores vol\u00famenes de ventas de acciones en cada sector industrial.<\/p>\n<p>Para dicha agregaci\u00f3n, se requerir\u00e1n varios pasos para traducir los datos a la forma necesaria (hablando en t\u00e9rminos generales).<\/p>\n<ol>\n<li>Crear una fuente basada en un tema que publique informaci\u00f3n sin procesar sobre el comercio de acciones. Tendremos que mapear un objeto de tipo StockTransaction a un objeto de tipo ShareVolume. El hecho es que el objeto StockTransaction contiene metadatos de ventas, y solo necesitamos los datos sobre la cantidad de acciones vendidas.<\/li>\n<li>Agrupar los datos de ShareVolume por s\u00edmbolos de acciones. Despu\u00e9s de agrupar por s\u00edmbolos, se pueden resumir estos datos a cantidades intermedias de volumen de ventas de acciones. Cabe se\u00f1alar que el m\u00e9todo KStream.groupBy devuelve una instancia de tipo KGroupedStream. Y para obtener una instancia de KTable, se puede invocar el m\u00e9todo KGroupedStream.reduce.<\/li>\n<\/ol>\n<p><\/p>\n<blockquote><p><b>\u00bfQu\u00e9 es la interfaz KGroupedStream?<\/b><\/p>\n<p>Los m\u00e9todos KStream.groupBy y KStream.groupByKey devuelven una instancia de KGroupedStream. KGroupedStream es una representaci\u00f3n intermedia de un flujo de eventos despu\u00e9s de agrupar por claves. No est\u00e1 destinado para ser manipulado directamente. En cambio, KGroupedStream se utiliza para operaciones de agregaci\u00f3n, cuyo resultado siempre es una KTable. Y dado que el resultado de las operaciones de agregaci\u00f3n es una KTable y se aplica un almacenamiento de estado, es posible que no todas las actualizaciones del resultado se env\u00eden a lo largo de la tuber\u00eda.<\/p>\n<p>El m\u00e9todo KTable.groupBy devuelve una KGroupedTable similar: una representaci\u00f3n intermedia de un flujo de actualizaciones, reorganizadas por clave.<\/p><\/blockquote>\n<p>\nTomemos un peque\u00f1o descanso y veamos la figura 5.9, que muestra lo que hemos logrado. Esta topolog\u00eda ya deber\u00eda ser familiar para usted.<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/9fd61317cde376362adcaeec72908919.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nAhora veamos el c\u00f3digo para esta topolog\u00eda (se puede encontrar en el archivo src\/main\/java\/bbejeck\/chapter_5\/AggregationsAndReducingExample.java) (listado 5.2).<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/f937287e448295fbd467c283ceca316a.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nEl c\u00f3digo proporcionado es conciso y realiza muchas acciones en unas pocas l\u00edneas. En el primer par\u00e1metro del m\u00e9todo builder.stream, puede notar algo nuevo para usted: el valor de la enumeraci\u00f3n AutoOffsetReset.EARLIEST (tambi\u00e9n existe LATEST), que se establece mediante el m\u00e9todo Consumed.withOffsetResetPolicy. Con esta enumeraci\u00f3n, puede especificar la estrategia de reinicio de desplazamientos para cada KStream o KTable, que tiene prioridad sobre la configuraci\u00f3n de reinicio de desplazamientos.<\/p>\n<blockquote><p><b>GroupByKey y GroupBy<\/b><\/p>\n<p>En la interfaz KStream hay dos m\u00e9todos para agrupar registros: GroupByKey y GroupBy. Ambos devuelven KGroupedTable, por lo que puede surgir la pregunta l\u00f3gica: \u00bfcu\u00e1l es la diferencia entre ellos y cu\u00e1ndo usar cada uno?<\/p>\n<p>El m\u00e9todo GroupByKey se aplica cuando las claves en KStream ya no est\u00e1n vac\u00edas. Lo m\u00e1s importante es que la bandera \"requiere re-particionamiento\" nunca se estableci\u00f3.<\/p>\n<p>El m\u00e9todo GroupBy supone que se han cambiado las claves para la agrupaci\u00f3n, por lo que la bandera de re-particionamiento se establece en verdadero. La ejecuci\u00f3n despu\u00e9s del m\u00e9todo GroupBy de uniones, agregaciones, etc., resultar\u00e1 en un re-particionamiento autom\u00e1tico.<br \/>\nResumen: se debe utilizar GroupByKey siempre que sea posible, en lugar de GroupBy.<\/p><\/blockquote>\n<p>\nLo que hacen los m\u00e9todos mapValues y groupBy es claro, as\u00ed que veamos el m\u00e9todo sum() (que se puede encontrar en el archivo src\/main\/java\/bbejeck\/model\/ShareVolume.java) (listado 5.3).<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/8e9a6f873594f9b5fef9a96c42353a61.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nEl m\u00e9todo ShareVolume.sum devuelve la suma intermedia del volumen de acciones vendidas, y el resultado de toda la cadena de c\u00e1lculos representa un objeto KTable. Ahora entiendes el papel que juega KTable. A medida que llegan objetos ShareVolume, la \u00faltima actualizaci\u00f3n relevante se conserva en el respectivo objeto KTable. Es importante recordar que todas las actualizaciones se reflejan en el anterior shareVolumeKTable, pero no todas se env\u00edan a continuaci\u00f3n.<\/p>\n<p>A continuaci\u00f3n, utilizando este KTable, realizamos la agregaci\u00f3n (por la cantidad de acciones vendidas) para obtener cinco empresas con los mayores vol\u00famenes de ventas de acciones en cada una de las industrias. Nuestras acciones ser\u00e1n similares a las del primer agrupamiento.<\/p>\n<ol>\n<li>Realizar otra operaci\u00f3n groupBy para agrupar objetos ShareVolume individuales por industria.<\/li>\n<li>Proceder a sumar objetos ShareVolume. Esta vez, el objeto de agregaci\u00f3n es una cola de prioridad de tama\u00f1o fijo. En esta cola de tama\u00f1o fijo se conservan \u00fanicamente cinco empresas con las mayores cantidades de acciones vendidas.<\/li>\n<li>Convertir las colas del punto anterior en un valor string y devolver cinco de las m\u00e1s vendidas por cantidad de acciones por sectores industriales.<\/li>\n<li>Registrar los resultados en formato string en el t\u00f3pico.<\/li>\n<\/ol>\n<p>\nEn la figura 5.10 se muestra un gr\u00e1fico de la topolog\u00eda del flujo de datos. Como se puede ver, el segundo ciclo de procesamiento es bastante simple.<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/dabd1507eee267038edb7f8d76d8d8ae.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nAhora, habiendo entendido claramente la estructura de este segundo ciclo de procesamiento, se puede consultar su c\u00f3digo fuente (lo encontrar\u00e1 en el archivo src\/main\/java\/bbejeck\/chapter_5\/AggregationsAndReducingExample.java) (listado 5.4).<\/p>\n<p>En este inicializador hay una variable fixedQueue. Este es un objeto personalizado: un adaptador para java.util.TreeSet, que se utiliza para rastrear los N mayores resultados en orden decreciente de la cantidad de acciones vendidas.<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/488b072b0d91b81c925ca72291e69e48.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nYa te has encontrado con las llamadas groupBy y mapValues, as\u00ed que no nos detendremos en ellas (llamamos al m\u00e9todo KTable.toStream, ya que el m\u00e9todo KTable.print se considera obsoleto). Sin embargo, a\u00fan no has visto la versi\u00f3n KTable del m\u00e9todo aggregate(), as\u00ed que dedicaremos un poco de tiempo a discutirlo.<\/p>\n<p>Como recordar\u00e1s, KTable se distingue en que los registros con las mismas claves se consideran actualizaciones. KTable reemplaza el registro antiguo por el nuevo. La agregaci\u00f3n ocurre de manera similar: se agregan los \u00faltimos registros con la misma clave. Cuando llega un registro, se agrega a la instancia de la clase FixedSizePriorityQueue utilizando un sumador (el segundo par\u00e1metro en la llamada al m\u00e9todo aggregate), pero si ya existe otro registro con la misma clave, el registro antiguo se elimina mediante un restador (el tercer par\u00e1metro en la llamada al m\u00e9todo aggregate).<\/p>\n<p>Todo esto significa que nuestro agregador, FixedSizePriorityQueue, no agrega todos los valores con una misma clave, sino que mantiene una suma m\u00f3vil de las cantidades de las N acciones m\u00e1s vendidas. Cada registro que llega contiene el total de acciones vendidas hasta ese momento. KTable te proporcionar\u00e1 informaci\u00f3n sobre qu\u00e9 acciones de qu\u00e9 empresas se est\u00e1n vendiendo m\u00e1s en este momento, no se requiere agregaci\u00f3n m\u00f3vil de cada actualizaci\u00f3n.<\/p>\n<p>Hemos aprendido a hacer dos cosas importantes:<\/p>\n<ul>\n<li>agrupar valores en KTable por una clave com\u00fan;<\/li>\n<li>realizar operaciones \u00fatiles sobre esos valores agrupados, como reducci\u00f3n y agregaci\u00f3n.<\/li>\n<\/ul>\n<p>\nLa habilidad para realizar estas operaciones es importante para comprender el significado de los datos que fluyen a trav\u00e9s de la aplicaci\u00f3n Kafka Streams y determinar qu\u00e9 informaci\u00f3n est\u00e1n transmitiendo.<\/p>\n<p>Tambi\u00e9n hemos reunido algunos de los conceptos clave discutidos previamente en este libro. En el cap\u00edtulo 4, hablamos de lo crucial que es para una aplicaci\u00f3n de streaming tener un estado local y tolerante a fallos. El primer ejemplo de este cap\u00edtulo demostr\u00f3 por qu\u00e9 el estado local es tan importante: permite hacer un seguimiento de la informaci\u00f3n que ya has visto. El acceso local evita retrasos en la red, lo que hace que la aplicaci\u00f3n sea m\u00e1s eficiente y resistente a errores.<\/p>\n<p>Al realizar cualquier operaci\u00f3n de reducci\u00f3n o agregaci\u00f3n, es necesario especificar el nombre del almac\u00e9n de estado. Las operaciones de reducci\u00f3n y agregaci\u00f3n devuelven una instancia de KTable, y KTable utiliza el almac\u00e9n de estado para reemplazar los antiguos resultados por nuevos. Como has visto, no todas las actualizaciones se env\u00edan m\u00e1s adelante en la canalizaci\u00f3n, y esto es importante, ya que las operaciones de agregaci\u00f3n est\u00e1n destinadas a obtener informaci\u00f3n final. Si no se aplica un estado local, KTable enviar\u00e1 todos los resultados de agregaci\u00f3n y reducci\u00f3n.<\/p>\n<p>A continuaci\u00f3n, veremos c\u00f3mo realizar operaciones como la agregaci\u00f3n dentro de un intervalo de tiempo espec\u00edfico, conocidas como operaciones de ventana (windowing operations).<\/p>\n<h3>5.3.2. Operaciones de Ventana<\/h3>\n<p>\nEn la secci\u00f3n anterior, nos familiarizamos con la reducci\u00f3n y la agregaci\u00f3n 'm\u00f3viles'. La aplicaci\u00f3n realizaba una reducci\u00f3n continua del volumen de ventas de acciones, seguida de la agregaci\u00f3n de las cinco acciones m\u00e1s vendidas en el mercado.<\/p>\n<p>A veces, estas reducciones y agregaciones continuas son necesarias. Y otras veces, se requieren operaciones solo sobre un intervalo de tiempo determinado. Por ejemplo, calcular cu\u00e1ntas transacciones de acciones de una empresa espec\u00edfica se realizaron en los \u00faltimos 10 minutos. O cu\u00e1ntos usuarios hicieron clic en un nuevo banner publicitario en los \u00faltimos 15 minutos. La aplicaci\u00f3n puede realizar tales operaciones repetidamente, pero con resultados que solo se refieren a intervalos de tiempo espec\u00edficos (ventanas de tiempo).<\/p>\n<h3>Conteo de transacciones burs\u00e1tiles por comprador<\/h3>\n<p>\nEn el siguiente ejemplo, nos ocuparemos de rastrear transacciones burs\u00e1tiles de varios traders: ya sean grandes organizaciones o astutos financistas independientes.<\/p>\n<p>Existen dos posibles razones para este tipo de seguimiento. Una de ellas es la necesidad de saber qu\u00e9 compran\/venden los l\u00edderes del mercado. Si estos grandes jugadores e inversores experimentados ven oportunidades, tiene sentido seguir su estrategia. La segunda raz\u00f3n es la intenci\u00f3n de detectar cualquier posible se\u00f1al de transacciones ilegales utilizando informaci\u00f3n privilegiada. Para ello, necesitar\u00e1s analizar la correlaci\u00f3n entre los grandes picos de ventas y los comunicados de prensa importantes.<\/p>\n<p>Este seguimiento consta de etapas como:<\/p>\n<ul>\n<li>crear un flujo para leer del tema stock-transactions;<\/li>\n<li>agrupar las entradas entrantes por el identificador del comprador y el s\u00edmbolo de la acci\u00f3n. La llamada al m\u00e9todo groupBy devuelve una instancia de la clase KGroupedStream;<\/li>\n<li>devolver mediante el m\u00e9todo KGroupedStream.windowedBy un flujo de datos restringido por una ventana temporal, lo que permite realizar agregaciones en la ventana. Dependiendo del tipo de ventana, se devuelve ya sea TimeWindowedKStream o SessionWindowedKStream;<\/li>\n<li>contar las transacciones para la operaci\u00f3n de agregaci\u00f3n. El flujo de datos de la ventana determina si se incluye una entrada espec\u00edfica en este conteo;<\/li>\n<li>registrar los resultados en el tema o imprimirlos en la consola durante el desarrollo.<\/li>\n<\/ul>\n<p>\nLa topolog\u00eda de esta aplicaci\u00f3n es simple, pero una representaci\u00f3n visual no estar\u00eda de m\u00e1s. Veamos la figura 5.11.<\/p>\n<p>A continuaci\u00f3n, examinaremos la funcionalidad de las operaciones de ventana y el c\u00f3digo correspondiente.<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/67d9d8d8acb1370a7c7fe5cd9a8b6aa7.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/p>\n<h3>Tipos de ventanas<\/h3>\n<p>\nEn Kafka Streams existen tres tipos de ventanas:<\/p>\n<ul>\n<li>de sesiones;<\/li>\n<li>de 'bote' (tumbling);<\/li>\n<li>deslizantes\/'saltadoras' (sliding\/hopping).<\/li>\n<\/ul>\n<p>\nCu\u00e1l elegir depende de los requisitos del negocio. Las ventanas 'bote' y 'saltadoras' se limitan por tiempo, mientras que las restricciones de las sesiones est\u00e1n relacionadas con las acciones de los usuarios: la duraci\u00f3n de la(s) sesi\u00f3n(es) se define exclusivamente por la actividad del usuario. Lo principal es recordar que todos los tipos de ventanas se basan en las marcas de fecha\/hora de las entradas, no en el tiempo del sistema.<\/p>\n<p>A continuaci\u00f3n, implementaremos nuestra topolog\u00eda con cada uno de los tipos de ventanas. El c\u00f3digo completo solo se proporcionar\u00e1 en el primer ejemplo, para los otros tipos de ventanas no habr\u00e1 cambios, excepto el tipo de operaci\u00f3n de ventana.<\/p>\n<h3>Ventanas de sesiones<\/h3>\n<p>\nLas ventanas de sesi\u00f3n son muy diferentes de todos los dem\u00e1s tipos de ventanas. Est\u00e1n limitadas no tanto por el tiempo, sino por la actividad del usuario (o la actividad de la entidad que desear\u00edas rastrear). Las ventanas de sesi\u00f3n se delimitan por periodos de inactividad.<\/p>\n<p>La figura 5.12 ilustra el concepto de ventanas de sesi\u00f3n. Una sesi\u00f3n m\u00e1s peque\u00f1a se mezclar\u00e1 con la sesi\u00f3n que est\u00e1 a su izquierda. Y la sesi\u00f3n a la derecha ser\u00e1 separada, ya que sigue a un largo periodo de inactividad. Las ventanas de sesi\u00f3n se basan en las acciones de los usuarios, pero aplican marcas de fecha\/hora de los registros para determinar a qu\u00e9 sesi\u00f3n pertenece el registro.<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/d39a5db7af7d6aa2194802622b2b47fd.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/p>\n<h3>Uso de ventanas de sesi\u00f3n para rastrear transacciones burs\u00e1tiles.<\/h3>\n<p>\nUtilizaremos ventanas de sesi\u00f3n para capturar informaci\u00f3n sobre transacciones burs\u00e1tiles. La implementaci\u00f3n de las ventanas de sesi\u00f3n se muestra en el listado 5.5 (que se puede encontrar en el archivo src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKTableJoinExample.java).<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/2dcbd9a36baec0e746aad165121451b3.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nLa mayor\u00eda de las operaciones de esta topolog\u00eda ya las has encontrado, as\u00ed que no es necesario revisarlas nuevamente aqu\u00ed. Pero hay algunos nuevos elementos que discutiremos ahora.<\/p>\n<p>En cada operaci\u00f3n groupBy generalmente se realiza alguna operaci\u00f3n de agregaci\u00f3n (agregaci\u00f3n, reducci\u00f3n o conteo). Se puede realizar agregaci\u00f3n acumulativa con un total acumulado, o agregaci\u00f3n por ventana, donde se consideran los registros dentro de un intervalo de tiempo espec\u00edfico.<\/p>\n<p>El c\u00f3digo del listado 5.5 cuenta la cantidad de transacciones dentro de las ventanas de sesi\u00f3n. En la figura 5.13, estas acciones se analizan paso a paso.<\/p>\n<p>Con la llamada windowedBy(SessionWindows.with(twentySeconds).until(fifteenMinutes)) estamos creando una ventana de sesi\u00f3n con un intervalo de inactividad de 20 segundos y un intervalo de retenci\u00f3n de 15 minutos. El intervalo de inactividad de 20 segundos significa que la aplicaci\u00f3n incluir\u00e1 cualquier registro que llegue dentro de los 20 segundos desde el final o inicio de la sesi\u00f3n actual en la sesi\u00f3n actual (activa).<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/9bd47b04698086872fd135b4c67eb938.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nA continuaci\u00f3n, indicamos qu\u00e9 operaci\u00f3n de agregaci\u00f3n se debe realizar en la ventana de sesi\u00f3n: en este caso, contar. Si el registro de entrada sale del intervalo de inactividad (por cualquiera de los lados de la marca de fecha\/hora), la aplicaci\u00f3n crea una nueva sesi\u00f3n. El intervalo de retenci\u00f3n significa mantener la sesi\u00f3n durante un tiempo determinado y permite datos tard\u00edos que salen del per\u00edodo de inactividad de la sesi\u00f3n, pero que a\u00fan pueden ser agregados. Adem\u00e1s, el inicio y el final de una nueva sesi\u00f3n, resultante de la fusi\u00f3n, corresponden a la marca de fecha\/hora m\u00e1s temprana y m\u00e1s tard\u00eda.<\/p>\n<p>Veamos algunos registros del m\u00e9todo count para entender c\u00f3mo funcionan las sesiones (Tabla 5.1).<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/ec04aae466d88c2d2349474069c8d541.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nAl recibir registros, buscamos sesiones existentes con la misma clave, donde el tiempo de finalizaci\u00f3n es menor que la marca de fecha\/hora actual \u2014 intervalo de inactividad \u2014 y el tiempo de inicio es mayor que la marca de fecha\/hora actual m\u00e1s el intervalo de inactividad. Teniendo en cuenta esto, cuatro registros de la Tabla 5.1 se fusionan en una sola sesi\u00f3n de la siguiente manera.<\/p>\n<p>1. Primero llega el registro 1, as\u00ed que el tiempo de inicio es igual al tiempo de finalizaci\u00f3n y es 00:00:00.<\/p>\n<p>2. Luego llega el registro 2, y buscamos sesiones que finalicen no antes de las 23:59:55 y que comiencen no m\u00e1s tarde de las 00:00:35. Encontramos el registro 1 y fusionamos las sesiones 1 y 2. Tomamos el tiempo de inicio de la sesi\u00f3n 1 (el m\u00e1s temprano) y el tiempo de finalizaci\u00f3n de la sesi\u00f3n 2 (el m\u00e1s tard\u00edo), por lo que nuestra nueva sesi\u00f3n comienza a las 00:00:00 y finaliza a las 00:00:15.<\/p>\n<p>3. Llega el registro 3, buscamos sesiones entre las 00:00:30 y las 00:01:10 y no encontramos ninguna. Agregamos una segunda sesi\u00f3n para la clave 123-345-654,FFBE, que comienza y termina a las 00:00:50.<\/p>\n<p>4. Llega el registro 4, y buscamos sesiones entre las 23:59:45 y las 00:00:25. Esta vez encontramos ambas sesiones: 1 y 2. Las tres sesiones se fusionan en una, con un tiempo de inicio de 00:00:00 y un tiempo de finalizaci\u00f3n de 00:00:15.<\/p>\n<p>De lo expuesto en esta secci\u00f3n, es importante recordar los siguientes matices:<\/p>\n<ul>\n<li>las sesiones no son ventanas de tama\u00f1o fijo. La duraci\u00f3n de la sesi\u00f3n se determina por la actividad dentro del intervalo de tiempo establecido;<\/li>\n<li>las marcas de fecha\/hora en los datos determinan si un evento entra en una sesi\u00f3n existente o en el intervalo de inactividad.<\/li>\n<\/ul>\n<p>\nA continuaci\u00f3n, discutiremos la siguiente variante de ventanas: las ventanas \"rodantes\".<\/p>\n<h3>Ventanas \"rodantes\"<\/h3>\n<p>\nLas ventanas \u00abrodantes\u00bb (tumbling) capturan eventos que ocurren en un intervalo de tiempo espec\u00edfico. Imagina que necesitas capturar todas las transacciones burs\u00e1tiles de una empresa cada 20 segundos, as\u00ed que recopilas todos los eventos durante ese intervalo. Al final del intervalo de 20 segundos, la ventana \u00abrueda\u00bb y pasa a un nuevo intervalo de observaci\u00f3n de 20 segundos. La Figura 5.14 ilustra esta situaci\u00f3n.<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/e90e560d9ddda5e2e524b7c387ad9874.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nComo puedes ver, todos los eventos que llegaron en los \u00faltimos 20 segundos est\u00e1n incluidos en la ventana. Al final de este per\u00edodo, se crea una nueva ventana.<\/p>\n<p>En la Listado 5.6, se muestra el c\u00f3digo que demuestra el uso de ventanas \u00abrodantes\u00bb para capturar cada 20 segundos las transacciones burs\u00e1tiles (se puede encontrar en el archivo src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java).<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/e055bb1b288c7d500b64372fe3fbf064.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nGracias a este peque\u00f1o cambio en la llamada al m\u00e9todo TimeWindows.of, se puede usar una ventana \u00abrodante\u00bb. En este ejemplo, no hay llamada al m\u00e9todo until(), por lo que se utilizar\u00e1 el intervalo de retenci\u00f3n predeterminado de 24 horas.<\/p>\n<p>Finalmente, es hora de pasar a la \u00faltima de las opciones de ventanas: las ventanas \u00abdeslizantes\u00bb (hopping).<\/p>\n<h3>Ventanas deslizantes (\u00abhopping\u00bb)<\/h3>\n<p>\nLas ventanas deslizantes\/\u00abhopping\u00bb (sliding\/hopping) son similares a las \u00abrodantes\u00bb, pero con una peque\u00f1a diferencia. Las ventanas deslizantes no esperan a que termine el intervalo de tiempo antes de crear una nueva ventana para procesar eventos recientes. Inician nuevos c\u00e1lculos despu\u00e9s de un intervalo de espera que es menor que la duraci\u00f3n de la ventana.<\/p>\n<p>Para ilustrar las diferencias entre las ventanas \u00abrodantes\u00bb y las \u00abdeslizantes\u00bb, volvamos al ejemplo de contar transacciones burs\u00e1tiles. Nuestro objetivo sigue siendo contar el n\u00famero de transacciones, pero no queremos esperar todo el intervalo de tiempo antes de actualizar el contador. En cambio, vamos a actualizar el contador a intervalos m\u00e1s cortos de tiempo. Por ejemplo, seguiremos contando el n\u00famero de transacciones cada 20 segundos, pero actualizaremos el contador cada 5 segundos, como se muestra en la Fig. 5.15. Esto nos da tres ventanas de resultados con datos superpuestos.<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/415c8cd9f2b60d453a1a01c3bc99331f.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nEn la Listado 5.7, se proporciona el c\u00f3digo para definir ventanas deslizantes (se puede encontrar en el archivo src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java).<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/ab2d1a64380d256d3fb084e16597417c.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nUna ventana \u00abrebotante\u00bb puede transformarse en una ventana \u00absaltadora\u00bb a\u00f1adiendo la llamada al m\u00e9todo advanceBy(). En el ejemplo proporcionado, el intervalo de retenci\u00f3n es de 15 minutos.<\/p>\n<p>En esta secci\u00f3n, hemos visto c\u00f3mo limitar los resultados de la agregaci\u00f3n mediante ventanas temporales. En particular, me gustar\u00eda que recordaras las siguientes tres cosas de esta secci\u00f3n:<\/p>\n<ul>\n<li>el tama\u00f1o de las ventanas de sesi\u00f3n no est\u00e1 limitado por el intervalo de tiempo, sino por la actividad de los usuarios;<\/li>\n<li>las ventanas \u00abrebotantes\u00bb brindan una visi\u00f3n de los eventos dentro de un periodo de tiempo determinado;<\/li>\n<li>la duraci\u00f3n de las ventanas \u00absaltadoras\u00bb es fija, pero se actualizan con frecuencia y pueden contener registros superpuestos en todas las ventanas.<\/li>\n<\/ul>\n<p>\nA continuaci\u00f3n, veremos c\u00f3mo convertir KTable de nuevo en KStream para realizar un join.<\/p>\n<h3>5.3.3. Join de objetos KStream y KTable<\/h3>\n<p>\nEn el cap\u00edtulo 4 discutimos c\u00f3mo unir dos objetos KStream. Ahora aprenderemos a unir KTable y KStream. Esto puede ser necesario por la siguiente raz\u00f3n sencilla. KStream es un flujo de registros, y KTable es un flujo de actualizaciones de registros, pero a veces puede ser necesario agregar un contexto adicional al flujo de registros mediante actualizaciones de KTable.<\/p>\n<p>Tomemos los datos sobre el n\u00famero de transacciones burs\u00e1tiles y los unamos con las noticias burs\u00e1tiles de las correspondientes industrias. Esto es lo que debemos hacer para lograrlo, considerando el c\u00f3digo ya existente.<\/p>\n<ol>\n<li>Convertir el objeto KTable con datos sobre el n\u00famero de transacciones burs\u00e1tiles en KStream, reemplazando la clave por una clave que indique la industria correspondiente a dicho s\u00edmbolo de acciones.<\/li>\n<li>Crear un objeto KTable que lea datos del tema con las noticias burs\u00e1tiles. Este nuevo KTable ser\u00e1 categorizado por industrias.<\/li>\n<li>Unir actualizaciones de noticias con informaci\u00f3n sobre el n\u00famero de transacciones burs\u00e1tiles por industria.<\/li>\n<\/ol>\n<p>\nAhora veamos c\u00f3mo implementar este plan de acci\u00f3n.<\/p>\n<h3>Conversi\u00f3n de KTable a KStream<\/h3>\n<p>\nPara convertir KTable a KStream, es necesario hacer lo siguiente.<\/p>\n<ol>\n<li>Llamar al m\u00e9todo KTable.toStream().<\/li>\n<li>Usando la llamada al m\u00e9todo KStream.map, reemplazar la clave con el nombre de la industria, despu\u00e9s extraer del objeto Windowed el objeto TransactionSummary.<\/li>\n<\/ol>\n<p>\nVincularemos estas operaciones en una cadena de la siguiente manera (el c\u00f3digo se puede encontrar en el archivo src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java) (listado 5.8).<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/0d43c2650f6e66e2816ed383da3a29c2.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nDado que estamos realizando la operaci\u00f3n KStream.map, el reagrupamiento para la instancia KStream devuelta se realiza autom\u00e1ticamente cuando se utiliza en una uni\u00f3n.<\/p>\n<p>Hemos completado el proceso de transformaci\u00f3n, ahora necesitamos crear un objeto KTable para leer las noticias de la bolsa.<\/p>\n<h3>Creaci\u00f3n de KTable para noticias de la bolsa<\/h3>\n<p>\nAfortunadamente, para crear un objeto KTable es suficiente una l\u00ednea de c\u00f3digo (este c\u00f3digo se puede encontrar en el archivo src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java) (listado 5.9).<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/6e83a393fdc9ab74fda4cbdddddb5213.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nCabe destacar que no es necesario especificar ning\u00fan objeto Serde, ya que se utilizan Serde de cadenas en la configuraci\u00f3n. Adem\u00e1s, gracias al uso de la enumeraci\u00f3n EARLIEST, la tabla se llena con registros desde el principio.<\/p>\n<p>Ahora podemos pasar al paso final: la uni\u00f3n.<\/p>\n<h3>Uni\u00f3n de actualizaciones de noticias con datos sobre el n\u00famero de transacciones<\/h3>\n<p>\nLa creaci\u00f3n de la uni\u00f3n no presenta dificultades. Usaremos una uni\u00f3n izquierda en caso de que no haya noticias de bolsa para la industria correspondiente (el c\u00f3digo necesario se puede encontrar en el archivo src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java) (listado 5.10).<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/92ed70f98927d2f778ad14dbc2a5aa26.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nEste operador leftJoin es bastante simple. A diferencia de las uniones del cap\u00edtulo 4, el m\u00e9todo JoinWindow no se utiliza, ya que al realizar la uni\u00f3n KStream-KTable hay solo un registro por clave en KTable. Esta uni\u00f3n no est\u00e1 restringida por el tiempo: un registro est\u00e1 presente en KTable o no. La conclusi\u00f3n principal es que los objetos KTable pueden enriquecer KStream con datos de referencia que se actualizan con menos frecuencia.<\/p>\n<p>Ahora consideraremos un m\u00e9todo m\u00e1s eficiente para enriquecer los eventos de KStream.<\/p>\n<h3>5.3.4. Objetos GlobalKTable<\/h3>\n<p>\nComo habr\u00e1s entendido, existe la necesidad de enriquecer los flujos de eventos o agregarles contexto. En el cap\u00edtulo 4 viste uniones de dos objetos KStream, y en la secci\u00f3n anterior, la uni\u00f3n de KStream y KTable. En todos estos casos, es necesario reagrupamiento del flujo de datos al mapear claves a un nuevo tipo o valor. A veces, el reagrupamiento se realiza expl\u00edcitamente, y otras veces Kafka Streams lo hace autom\u00e1ticamente. El reagrupamiento es necesario porque las claves han cambiado y los registros deben estar en nuevas secciones; de lo contrario, la uni\u00f3n ser\u00e1 imposible (esto se discuti\u00f3 en el cap\u00edtulo 4, en el apartado \"Reagrupamiento de datos\" de la subsecci\u00f3n 4.2.4).<\/p>\n<h3>La re-segmentaci\u00f3n tiene su precio<\/h3>\n<p>\nLa re-segmentaci\u00f3n requiere costos: costos adicionales de recursos para crear temas intermedios, mantener datos duplicados en otro tema; tambi\u00e9n significa un aumento en la latencia debido a la escritura y lectura de este tema. Adem\u00e1s, si es necesario realizar una uni\u00f3n en m\u00e1s de un aspecto o dimensi\u00f3n, se deben organizar las uniones en cadena, reflejar los registros con nuevas claves y volver a realizar el proceso de re-segmentaci\u00f3n.<\/p>\n<h3>Uniones con conjuntos de datos m\u00e1s peque\u00f1os<\/h3>\n<p>\nEn algunos casos, el volumen de datos de referencia con los que se planea hacer la uni\u00f3n es relativamente peque\u00f1o, por lo que las copias completas pueden caber localmente en cada uno de los nodos. Para tales situaciones, Kafka Streams proporciona la clase GlobalKTable.<\/p>\n<p>Las instancias de GlobalKTable son \u00fanicas, ya que la aplicaci\u00f3n replica todos los datos en cada uno de los nodos. Y dado que todos los nodos tienen todos los datos, no es necesario segmentar el flujo de eventos por la clave de los datos de referencia para que est\u00e9 disponible para todas las secciones. A trav\u00e9s de los objetos GlobalKTable tambi\u00e9n se pueden realizar uniones sin clave. Regresamos a uno de los ejemplos anteriores para demostrar esta capacidad.<\/p>\n<h3>Uni\u00f3n de objetos KStream con objetos GlobalKTable<\/h3>\n<p>\nEn la subsecci\u00f3n 5.3.2 realizamos una agregaci\u00f3n en ventana de transacciones burs\u00e1tiles por clientes. Los resultados de esta agregaci\u00f3n se ve\u00edan aproximadamente as\u00ed:<\/p>\n<pre><code class=\"plaintext\">{customerId='074-09-3705', stockTicker='GUTM'}, 17\n{customerId='037-34-5184', stockTicker='CORK'}, 16<\/code><\/pre>\n<p>\nAunque estos resultados cumpl\u00edan con el objetivo establecido, ser\u00eda m\u00e1s conveniente si tambi\u00e9n se mostrara el nombre del cliente y el nombre completo de la empresa. Para agregar el nombre del comprador y el nombre de la empresa, se pueden realizar uniones normales, pero eso requerir\u00eda hacer dos mapeos de claves y re-segmentaci\u00f3n. Con GlobalKTable se pueden evitar los costos de tales operaciones.<\/p>\n<p>Para ello, utilizaremos el objeto countStream del listado 5.11 (el c\u00f3digo correspondiente se puede encontrar en el archivo src\/main\/java\/bbejeck\/chapter_5\/GlobalKTableExample.java), uni\u00e9ndolo con dos objetos GlobalKTable.<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/fc4d91bbe062ceb94f5650224840b81e.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nYa hemos discutido esto antes, as\u00ed que no me repetir\u00e9. Pero se\u00f1alar\u00e9 que el c\u00f3digo en la funci\u00f3n toStream().map est\u00e1 abstra\u00eddo en un objeto-funci\u00f3n para mejorar la legibilidad, en lugar de usar una expresi\u00f3n lambda incrustada.<\/p>\n<p>El siguiente paso es declarar dos instancias de GlobalKTable (el c\u00f3digo correspondiente se puede encontrar en el archivo src\/main\/java\/bbejeck\/chapter_5\/GlobalKTableExample.java) (listado 5.12).<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/db3918895f174c8cfb5fc927b55f87a1.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/p>\n<p>Tenga en cuenta que los nombres de los temas se describen utilizando tipos enumerados.<\/p>\n<p>Ahora que hemos preparado todos los componentes, solo queda escribir el c\u00f3digo para la conexi\u00f3n (que se puede encontrar en el archivo src\/main\/java\/bbejeck\/chapter_5\/GlobalKTableExample.java) (listado 5.13).<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/799360cc99f1920c190a61fd4685d4ff.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nAunque hay dos uniones en este c\u00f3digo, est\u00e1n organizadas en forma de cadena, ya que por separado no se utiliza ninguno de sus resultados. Los resultados se muestran al final de toda la operaci\u00f3n.<\/p>\n<p>Al ejecutar la operaci\u00f3n de uni\u00f3n mencionada, obtendr\u00e1 resultados de la siguiente forma:<\/p>\n<pre><code class=\"plaintext\">{customer='Barney, Smith' company=\"Exxon\", transactions= 17}<\/code><\/pre>\n<p>\nLa esencia no ha cambiado, pero estos resultados son m\u00e1s claros.<\/p>\n<p>Si contamos el cap\u00edtulo 4, ya ha visto varios tipos de uniones en acci\u00f3n. Se enumeran en la tabla 5.2. Esta tabla refleja las capacidades de uni\u00f3n relevantes para la versi\u00f3n 1.0.0 de Kafka Streams; en futuras versiones, puede que algo cambie.<\/p>\n<p><img decoding=\"async\" alt=\"El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb\" src=\"\/wp-content\/uploads\/8e4cf35c64a8431bda43a5e748de275f.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nEn conclusi\u00f3n, recordar\u00e9 lo b\u00e1sico: puede unir flujos de eventos (KStream) y flujos de actualizaciones (KTable) utilizando un estado local. Adem\u00e1s, si el tama\u00f1o de los datos de referencia no es demasiado grande, puede utilizar el objeto GlobalKTable. GlobalKTable replica todas las particiones en cada uno de los nodos de la aplicaci\u00f3n Kafka Streams, lo que garantiza la disponibilidad de todos los datos sin importar qu\u00e9 partici\u00f3n corresponda a la clave.<\/p>\n<p>A continuaci\u00f3n, veremos una capacidad de Kafka Streams que permite observar cambios en el estado sin consumir datos del tema de Kafka.<\/p>\n<h3>5.3.5. Estado disponible para consultas<\/h3>\n<p>\nYa hemos realizado varias operaciones que involucran estado y siempre hemos mostrado los resultados en la consola (con fines de desarrollo) o los hemos grabado en un tema (con fines de producci\u00f3n). Al grabar resultados en un tema, es necesario usar un consumidor de Kafka para verlos.<\/p>\n<p>Leer datos de estos temas puede considerarse una forma de vistas materializadas. Para nuestros prop\u00f3sitos, podemos usar la definici\u00f3n de vista materializada de \"Wikipedia\": \"...un objeto f\u00edsico en la base de datos que contiene los resultados de la ejecuci\u00f3n de una consulta. Por ejemplo, puede ser una copia local de datos remotos, o un subconjunto de filas y\/o columnas de una tabla o de los resultados de un join, o una tabla agregada obtenida mediante agregaci\u00f3n\" (https:\/\/en.wikipedia.org\/wiki\/Materialized_view).<\/p>\n<p>Kafka Streams tambi\u00e9n permite realizar consultas interactivas a los almacenes de estado, lo que da la posibilidad de leer directamente estas vistas materializadas. Es importante destacar que la consulta al almac\u00e9n de estado es una operaci\u00f3n de 'solo lectura'. Gracias a esto, no tienes que preocuparte por hacer que el estado sea inconsistente durante el procesamiento de datos por la aplicaci\u00f3n.<\/p>\n<p>La posibilidad de realizar consultas directas a los almacenes de estado es de gran importancia. Significa que se pueden crear aplicaciones, como paneles de control, sin necesidad de obtener primero los datos del consumidor de Kafka. Tambi\u00e9n aumenta la eficiencia de la aplicaci\u00f3n, ya que no es necesario volver a escribir los datos:<\/p>\n<ul>\n<li>gracias a la localizaci\u00f3n de los datos, se puede acceder a ellos r\u00e1pidamente;<\/li>\n<li>se elimina la duplicaci\u00f3n de datos, ya que no se almacenan en un almac\u00e9n externo.<\/li>\n<\/ul>\n<p>\nLo principal que me gustar\u00eda que recordaras: puedes realizar consultas directamente sobre el estado desde la aplicaci\u00f3n. No se puede subestimar las posibilidades que esto te brinda. En lugar de consumir datos de Kafka y almacenar registros en la base de datos para la aplicaci\u00f3n, se pueden realizar consultas a los almacenes de estado con el mismo resultado. Las consultas directas a los almacenes de estado significan menos c\u00f3digo (sin consumidor) y menos software (sin necesidad de una tabla de base de datos para almacenar los resultados).<\/p>\n<p>Hemos cubierto una cantidad considerable de informaci\u00f3n en este cap\u00edtulo, as\u00ed que por ahora dejaremos de lado nuestra discusi\u00f3n sobre consultas interactivas a los almacenes de estado. Pero no se preocupen: en el cap\u00edtulo 9 crearemos una aplicaci\u00f3n sencilla: un panel informativo con consultas interactivas. Para demostrar las consultas interactivas y las posibilidades de su incorporaci\u00f3n en aplicaciones de Kafka Streams, utilizaremos algunos de los ejemplos de este y del cap\u00edtulo anterior.<\/p>\n<h3>Curr\u00edculum<\/h3>\n<p><\/p>\n<ul>\n<li>Los objetos KStream representan flujos de eventos, similares a inserciones en una base de datos. Los objetos KTable representan flujos de actualizaciones, siendo m\u00e1s parecidos a las actualizaciones en una base de datos. El tama\u00f1o de un objeto KTable no crece, las entradas antiguas son reemplazadas por las nuevas.<\/li>\n<li>Los objetos KTable son necesarios para las operaciones de agregaci\u00f3n.<\/li>\n<li>Con las operaciones de ventana se pueden segmentar los datos agregados en cubos temporales.<\/li>\n<li>Gracias a los objetos GlobalKTable se puede acceder a los datos de referencia en cualquier parte de la aplicaci\u00f3n, independientemente de la partici\u00f3n por secciones.<\/li>\n<li>Es posible realizar uniones entre objetos KStream, KTable y GlobalKTable.<\/li>\n<\/ul>\n<p>\nHasta ahora nos hemos centrado en crear aplicaciones de Kafka Streams utilizando el DSL KStream de alto nivel. Aunque el enfoque de alto nivel permite escribir programas ordenados y concisos, su uso implica un cierto compromiso. Trabajar con el DSL KStream significa aumentar la concisi\u00f3n del c\u00f3digo a costa de reducir el control. En el siguiente cap\u00edtulo abordaremos la API de bajo nivel de los nodos procesadores y exploraremos otros compromisos. Los programas ser\u00e1n m\u00e1s largos que los que hemos visto hasta ahora, pero tendremos la capacidad de crear pr\u00e1cticamente cualquier nodo procesador que podamos necesitar.<\/p>\n<p>\u2192 Para m\u00e1s detalles, se puede consultar el libro en <noindex><a rel=\"nofollow\" href=\"https:\/\/www.piter.com\/collection\/best\/product\/kafka-streams-v-deystvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni\">el sitio del editor<\/a><\/noindex><\/p>\n<p>\u2192 Para los habitantes de Habr, un 25% de descuento con el cup\u00f3n \u2014 <b>Kafka Streams<\/b><\/p>\n<p>\u2192 Tras el pago de la versi\u00f3n impresa del libro, se enviar\u00e1 un libro electr\u00f3nico por correo electr\u00f3nico.<br \/>\n<br \/>Fuente: <a content=\"nofollow\" rel=\"nofollow\" href=\"https:\/\/habr.com\/ru\/company\/piter\/blog\/457756\/\">habr.com<\/a><\/p>","protected":false,"gt_translate_keys":[{"key":"rendered","format":"html"}]},"excerpt":{"rendered":"<p>\u041f\u0440\u0438\u0432\u0435\u0442, \u0425\u0430\u0431\u0440\u043e\u0436\u0438\u0442\u0435\u043b\u0438! \u042d\u0442\u0430 \u043a\u043d\u0438\u0433\u0430 \u043f\u043e\u0434\u043e\u0439\u0434\u0435\u0442 \u0434\u043b\u044f \u043b\u044e\u0431\u043e\u0433\u043e \u0440\u0430\u0437\u0440\u0430\u0431\u043e\u0442\u0447\u0438\u043a\u0430, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u0445\u043e\u0447\u0435\u0442 \u0440\u0430\u0437\u043e\u0431\u0440\u0430\u0442\u044c\u0441\u044f \u0432 \u043f\u043e\u0442\u043e\u043a\u043e\u0432\u043e\u0439 \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0435. \u041f\u043e\u043d\u0438\u043c\u0430\u043d\u0438\u0435 \u0440\u0430\u0441\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u043d\u043e\u0433\u043e \u043f\u0440\u043e\u0433\u0440\u0430\u043c\u043c\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u043f\u043e\u043c\u043e\u0436\u0435\u0442 \u043b\u0443\u0447\u0448\u0435 \u0438\u0437\u0443\u0447\u0438\u0442\u044c Kafka \u0438 Kafka Streams. \u0411\u044b\u043b\u043e \u0431\u044b \u043d\u0435\u043f\u043b\u043e\u0445\u043e \u0437\u043d\u0430\u0442\u044c \u0438 \u0441\u0430\u043c \u0444\u0440\u0435\u0439\u043c\u0432\u043e\u0440\u043a Kafka, \u043d\u043e \u044d\u0442\u043e \u043d\u0435 \u043e\u0431\u044f\u0437\u0430\u0442\u0435\u043b\u044c\u043d\u043e: \u044f \u0440\u0430\u0441\u0441\u043a\u0430\u0436\u0443 \u0432\u0430\u043c \u0432\u0441\u0435, \u0447\u0442\u043e \u043d\u0443\u0436\u043d\u043e. \u041e\u043f\u044b\u0442\u043d\u044b\u0435 \u0440\u0430\u0437\u0440\u0430\u0431\u043e\u0442\u0447\u0438\u043a\u0438 Kafka, \u043a\u0430\u043a \u0438 \u043d\u043e\u0432\u0438\u0447\u043a\u0438, \u0431\u043b\u0430\u0433\u043e\u0434\u0430\u0440\u044f \u044d\u0442\u043e\u0439 \u043a\u043d\u0438\u0433\u0435 \u043e\u0441\u0432\u043e\u044f\u0442 \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u0435 \u0438\u043d\u0442\u0435\u0440\u0435\u0441\u043d\u044b\u0445 \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0439 [&hellip;]<\/p>\n","protected":false,"gt_translate_keys":[{"key":"rendered","format":"html"}]},"author":1,"featured_media":0,"comment_status":"open","ping_status":"open","sticky":false,"template":"","format":"standard","meta":{"footnotes":""},"categories":[688],"tags":[],"class_list":["post-35786","post","type-post","status-publish","format-standard","hentry","category-administrirovanie"],"aioseo_notices":[],"aioseo_head":"\n\t\t<!-- All in One SEO 5.0.2 - aioseo.com -->\n\t<meta name=\"robots\" content=\"max-image-preview:large\" \/>\n\t<meta name=\"author\" content=\"Yuri Gagarin\"\/>\n\t<link rel=\"canonical\" href=\"https:\/\/prohoster.info\/es\/blog\/administrirovanie\/kniga-kafka-streams-v-dejstvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni\" \/>\n\t<meta name=\"generator\" content=\"All in One SEO (AIOSEO) 5.0.2\" \/>\n\t\t<meta property=\"og:locale\" content=\"es_ES\" \/>\n\t\t<meta property=\"og:site_name\" content=\"ProHoster | \u041a\u0443\u043f\u0438\u0442\u044c \u043d\u0430\u0434\u0435\u0436\u043d\u044b\u0439 \u0445\u043e\u0441\u0442\u0438\u043d\u0433 \u0434\u043b\u044f \u0441\u0430\u0439\u0442\u043e\u0432 \u0441 \u0437\u0430\u0449\u0438\u0442\u043e\u0439 \u043e\u0442 DDoS, VPS VDS \u0441\u0435\u0440\u0432\u0435\u0440\u044b\" \/>\n\t\t<meta property=\"og:type\" content=\"article\" \/>\n\t\t<meta property=\"og:title\" content=\"\ud83e\udd47\u041a\u043d\u0438\u0433\u0430 \u00abKafka Streams \u0432 \u0434\u0435\u0439\u0441\u0442\u0432\u0438\u0438. \u041f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u044f \u0438 \u043c\u0438\u043a\u0440\u043e\u0441\u0435\u0440\u0432\u0438\u0441\u044b \u0434\u043b\u044f \u0440\u0430\u0431\u043e\u0442\u044b \u0432 \u0440\u0435\u0430\u043b\u044c\u043d\u043e\u043c \u0432\u0440\u0435\u043c\u0435\u043d\u0438\u00bb | ProHoster\" \/>\n\t\t<meta property=\"og:url\" content=\"https:\/\/prohoster.info\/es\/blog\/administrirovanie\/kniga-kafka-streams-v-dejstvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni\" \/>\n\t\t<meta property=\"og:image\" content=\"https:\/\/prohoster.info\/wp-content\/uploads\/2021\/11\/logo-350.jpg\" \/>\n\t\t<meta property=\"og:image:secure_url\" content=\"https:\/\/prohoster.info\/wp-content\/uploads\/2021\/11\/logo-350.jpg\" \/>\n\t\t<meta property=\"og:image:width\" content=\"350\" \/>\n\t\t<meta property=\"og:image:height\" content=\"350\" \/>\n\t\t<meta property=\"article:published_time\" content=\"2019-10-31T19:06:19+00:00\" \/>\n\t\t<meta property=\"article:modified_time\" content=\"2019-10-31T19:06:19+00:00\" \/>\n\t\t<meta property=\"article:publisher\" content=\"https:\/\/www.facebook.com\/prohoster\" \/>\n\t\t<meta property=\"article:author\" content=\"https:\/\/www.facebook.com\/prohoster\" \/>\n\t\t<!-- All in One SEO -->\n\n","aioseo_head_json":{"title":"\ud83e\udd47El libro \u00abKafka Streams en acci\u00f3n. Aplicaciones y microservicios para trabajar en tiempo real\u00bb | ProHoster","description":"","canonical_url":"https:\/\/prohoster.info\/es\/blog\/administrirovanie\/kniga-kafka-streams-v-dejstvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni","robots":"max-image-preview:large","keywords":"","webmasterTools":{"miscellaneous":""},"schema":null,"og:locale":"es_ES","og:site_name":"ProHoster | \u041a\u0443\u043f\u0438\u0442\u044c \u043d\u0430\u0434\u0435\u0436\u043d\u044b\u0439 \u0445\u043e\u0441\u0442\u0438\u043d\u0433 \u0434\u043b\u044f \u0441\u0430\u0439\u0442\u043e\u0432 \u0441 \u0437\u0430\u0449\u0438\u0442\u043e\u0439 \u043e\u0442 DDoS, VPS VDS \u0441\u0435\u0440\u0432\u0435\u0440\u044b","og:type":"article","og:title":"\ud83e\udd47\u041a\u043d\u0438\u0433\u0430 \u00abKafka Streams \u0432 \u0434\u0435\u0439\u0441\u0442\u0432\u0438\u0438. \u041f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u044f \u0438 \u043c\u0438\u043a\u0440\u043e\u0441\u0435\u0440\u0432\u0438\u0441\u044b \u0434\u043b\u044f \u0440\u0430\u0431\u043e\u0442\u044b \u0432 \u0440\u0435\u0430\u043b\u044c\u043d\u043e\u043c \u0432\u0440\u0435\u043c\u0435\u043d\u0438\u00bb | ProHoster","og:url":"https:\/\/prohoster.info\/es\/blog\/administrirovanie\/kniga-kafka-streams-v-dejstvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni","og:image":"https:\/\/prohoster.info\/wp-content\/uploads\/2021\/11\/logo-350.jpg","og:image:secure_url":"https:\/\/prohoster.info\/wp-content\/uploads\/2021\/11\/logo-350.jpg","og:image:width":350,"og:image:height":350,"article:published_time":"2019-10-31T19:06:19+00:00","article:modified_time":"2019-10-31T19:06:19+00:00","article:publisher":"https:\/\/www.facebook.com\/prohoster","article:author":"https:\/\/www.facebook.com\/prohoster"},"aioseo_meta_data":{"post_id":"35786","title":null,"description":null,"keywords":null,"keyphrases":null,"primary_term":null,"canonical_url":null,"og_title":null,"og_description":null,"og_object_type":"default","og_image_type":"default","og_image_url":null,"og_image_width":null,"og_image_height":null,"og_image_custom_url":null,"og_image_custom_fields":null,"og_video":null,"og_custom_url":null,"og_article_section":null,"og_article_tags":null,"twitter_use_og":false,"twitter_card":"default","twitter_image_type":"default","twitter_image_url":null,"twitter_image_custom_url":null,"twitter_image_custom_fields":null,"twitter_title":null,"twitter_description":null,"schema":{"blockGraphs":[],"customGraphs":[],"default":{"data":{"Article":[],"Course":[],"Dataset":[],"FAQPage":[],"Movie":[],"Person":[],"Product":[],"ProductReview":[],"Car":[],"Recipe":[],"Service":[],"SoftwareApplication":[],"WebPage":[]},"graphName":"","isEnabled":true},"graphs":[]},"schema_type":null,"schema_type_options":null,"pillar_content":false,"robots_default":true,"robots_noindex":false,"robots_noarchive":false,"robots_nosnippet":false,"robots_nofollow":false,"robots_noimageindex":false,"robots_noodp":false,"robots_notranslate":false,"robots_max_snippet":null,"robots_max_videopreview":null,"robots_max_imagepreview":"large","priority":null,"frequency":null,"local_seo":null,"seo_analyzer_scan_date":"2026-01-22 00:45:19","breadcrumb_settings":null,"limit_modified_date":false,"reviewed_by":null,"ai":null,"created":"2021-03-01 01:56:32","updated":"2026-01-22 00:45:19","focus_keyword":null,"additional_keywords":null,"truseo_locale":null},"gt_translate_keys":[{"key":"link","format":"url"}],"_links":{"self":[{"href":"https:\/\/prohoster.info\/es\/wp-json\/wp\/v2\/posts\/35786","targetHints":{"allow":["GET"]}}],"collection":[{"href":"https:\/\/prohoster.info\/es\/wp-json\/wp\/v2\/posts"}],"about":[{"href":"https:\/\/prohoster.info\/es\/wp-json\/wp\/v2\/types\/post"}],"author":[{"embeddable":true,"href":"https:\/\/prohoster.info\/es\/wp-json\/wp\/v2\/users\/1"}],"replies":[{"embeddable":true,"href":"https:\/\/prohoster.info\/es\/wp-json\/wp\/v2\/comments?post=35786"}],"version-history":[{"count":0,"href":"https:\/\/prohoster.info\/es\/wp-json\/wp\/v2\/posts\/35786\/revisions"}],"wp:attachment":[{"href":"https:\/\/prohoster.info\/es\/wp-json\/wp\/v2\/media?parent=35786"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"https:\/\/prohoster.info\/es\/wp-json\/wp\/v2\/categories?post=35786"},{"taxonomy":"post_tag","embeddable":true,"href":"https:\/\/prohoster.info\/es\/wp-json\/wp\/v2\/tags?post=35786"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}