¿Es posible obtener el resultado final de la ventana en Kafka Streams al suprimir los resultados intermedios?
No puedo lograr este objetivo. ¿Qué está mal con mi código?
val builder = StreamsBuilder() builder.stream<String,Double>(inputTopic) .groupByKey() .windowedBy(TimeWindows.of(Duration.ofSeconds(15))) .count() .suppress(Suppressed.untilWindowCloses(unbounded())) // not working) .toStream() .print(Printed.toSysOut())Conduce a este error:
Failed to flush state store KSTREAM-AGGREGATE-STATE-STORE-0000000001: java.lang.ClassCastException: org.apache.kafka.streams.kstream.Windowed cannot be cast to java.lang.StringDetalles del código/error: https://gist.github.com/robie2011/1caa4772b60b5a6f993e6f98e792a380
El problema es una asimetría confusa en la forma en que Streams envuelve automáticamente los serdos explícitos durante la ventana, pero no envuelve automáticamente el serde predeterminado. En mi humilde opinión, este es un descuido que debe corregirse, por lo que he archivado: https://issues.apache.org/jira/browse/KAFKA-7806
Como han señalado otros, la solución es configurar explícitamente la clave serde en sentido ascendente y no confiar en la clave predeterminada serde. Tu también puedes:
Establezca los serdes en la agregación de ventanas con Materialized
val builder = StreamsBuilder() builder.stream<String,Double>(inputTopic) .groupByKey() .windowedBy(TimeWindows.of(Duration.ofSeconds(15))) .count(Materialized.with(Serdes.String(), Serdes.Long())) .suppress(Suppressed.untilWindowCloses(unbounded()))) .toStream() .print(Printed.toSysOut())(como recomendó Nishu)
(tenga en cuenta que no es necesario nombrar la operación de count , que tiene el efecto secundario de hacerla consultable)
O configure los serdes más arriba, por ejemplo en la entrada:
val builder = StreamsBuilder() builder.stream<String,Double>(inputTopic, Consumed.with(Serdes.String(), Serdes.Double())) .groupByKey() .windowedBy(TimeWindows.of(Duration.ofSeconds(15))) .count() .suppress(Suppressed.untilWindowCloses(unbounded()))) .toStream() .print(Printed.toSysOut())(como recomendó wardziniak)
La decisión es tuya; Creo que en este caso no es muy diferente en cualquier caso. Si estuviera haciendo una agregación diferente a count , probablemente estaría configurando el valor serde a través Materialized de todos modos, por lo que tal vez el primero sería un estilo más uniforme.
También noté que su definición de ventana no tiene un período de gracia establecido. El tiempo de cierre de la ventana se define como window end + grace period , y el valor predeterminado es de 24 horas, por lo que no vería nada emitido por la supresión hasta que se hayan ejecutado 24 horas de datos en la aplicación.
Para su esfuerzo de prueba, recomendaría probar:
.windowedBy(TimeWindows.of(Duration.ofSeconds(15)).grace(Duration.ZERO))En producción, querrá seleccionar un período de gracia que equilibre la cantidad de retraso de eventos que espera en su transmisión con la cantidad de prontitud de emisión que desea ver de la supresión.
Una nota final, noté en su esencia que no ha cambiado el almacenamiento en caché predeterminado o el intervalo de confirmación. Como resultado, notará que el propio operador de count almacenará en búfer las actualizaciones durante los 30 segundos predeterminados antes de pasarlas a la supresión. Esta es una buena configuración para la producción, por lo que no crea un cuello de botella en su disco local o en el agente de Kafka. Pero puede que te sorprenda mientras estás probando.
Por lo general, para las pruebas (o probar cosas de forma interactiva), deshabilitaré el almacenamiento en caché y configuraré el intervalo de compromiso corto para la máxima cordura del desarrollador:
properties.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, 0); properties.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 100);Perdón por el descuido serde. Espero que abordemos KAFKA-7806 pronto.
¡Espero que esto ayude!
El problema es con el KeySerde. Dado que la operación WindowedBy da como resultado una clave de tipo Windowed<String> , pero .suppress() está usando un tipo de clave predeterminado.
Por lo tanto, debe definir KeySerde en la tienda State mientras llama al método de conteo como se indica a continuación:
builder.stream<String,Double>inputTopic) .groupByKey() .windowedBy(TimeWindows.of(Duration.ofSeconds(15))) .count(Materialized.<String, Long, WindowStore<Bytes,byte[]>>as("count").withCachingDisabled().withKeySerde(Serdes.String())) .suppress(Suppressed.untilWindowCloses(BufferConfig.unbounded())) .toStream() . print(Printed.toSysOut());Agregue Consumed , cuando cree una secuencia: builder.stream<String,Double>(inputTopic, Consumed. with (Serdes.String(), Serdes.Double())