Tengo un productor de flujo que emitirá elementos con un período aleatorio, no quiero manejar estos elementos una vez emitidos, prefiero reunirlos en una lista, luego manejarlos y mantener este proceso en ejecución.
Por ejemplo, si quiero tener 2s como un período, lo que quiero obtener para el siguiente flujo es
val flow = flow{ for(i in 1..5){ emit(i) delay(1100) } } //#1 handle [1,2] //#2 handle [3,4] //#3 handle [5]A continuación se muestra lo que intenté pero fallé:
take(2) , esto hará que el flujo finalice después de obtener dos elementos, y no hay tiempo para controlarcollectToList().sample(2s) , esta es mi función personalizada, agregando elementos a una lista. Pero el problema es que necesito borrar la lista en collectToList() después de sample() , de lo contrario, la lista se agregará para siempre, pero si hago el borrado, esto puede causar problemas de concurrencia como "agregar un nuevo elemento a la lista solo antes de borrar la lista"¿Alguna idea mejor?
La respuesta corta es que no hay una función integrada para esto en este momento AFAIK.
Creo que lo que está buscando es un operador chunked basado en el tiempo, que se ha discutido bastante , pero aún no ha llegado a la biblioteca. No estoy seguro si está en sus planes para el futuro cercano.
Tal vez pueda echar un vistazo a la implementación propuesta en esta solicitud de combinación y adaptarla a sus necesidades (probablemente no necesite tanta flexibilidad, ni las estrategias de fragmentación adicionales). Lo siento, estoy afk, así que no puedo hacer la adaptación yo mismo.