¿Es seguro seguir el código y por qué?
val flow: Flow<String> = ... val allStrings = mutableListOf<String>() var sum = 0 flow.transform { allStrings += it emit(it.toInt()) }.collect { sum += it } La siguiente prueba demuestra que se llama al cuerpo de collect {} desde diferentes subprocesos:
val ctx = newFixedThreadPoolContext(32, "my-context") runBlocking(ctx) { val f = flow<Int> { (1 .. 1000).forEach { emit(it) } } var t: Thread? = null f.collect { delay(1) // this requirement will fail if (t == null) t = Thread.currentThread() else require(t == Thread.currentThread()) } }Y otro que pone a prueba la publicación:
fun main(args: Array<String>) { val ctx = newFixedThreadPoolContext(32, "my-context") runBlocking(ctx) { val f = flow<Int> { (1 .. 1000000).forEach { emit(it) } } var c = 0 f.transform { c += 1 boo() c += 1 emit(it) c += 1 }.collect { c += 1 boo() c += 1 } println(c) // prints 5_000_000 } } suspend fun boo() { withContext(Dispatchers.IO) { } }Por lo tanto, parece que el flujo de kotlin garantiza la publicación entre invocaciones de rutina, pero ¿es este un efecto secundario intencional (o incluso documentado) o de implementación?
Sí, el código es seguro. Los flujos son secuenciales por defecto.
Cada colección individual de un flujo se realiza secuencialmente a menos que se utilicen operadores especiales que operan en múltiples flujos. El cobro funciona directamente en la corrutina que llama a un operador de terminal (
collect). No se lanzan nuevas rutinas por defecto. Cada valor emitido es procesado por todos los operadores intermedios desde aguas arriba hasta aguas abajo y luego se entrega al operador terminal después.
Los flujos de Kotlin se basan en funciones de suspensión y son completamente secuenciales, pero no de un solo subproceso.
Por lo tanto, parece que el flujo de kotlin garantiza la publicación entre invocaciones de rutina, pero ¿es este un efecto secundario intencional (o incluso documentado) o de implementación?
un flujo es un tipo que puede emitir múltiples valores secuencialmente
De acuerdo con esto, entiendo que no se recopilará ningún valor nuevo hasta que se procese el valor anterior, sin importar cuántos subprocesos estén involucrados.
Como menciona Roman en su comentario, aquí hay un buen artículo sobre la ejecución secuencial . Una gran cita de allí:
Aunque una corrutina en Kotlin puede ejecutarse en varios subprocesos, es como un subproceso desde el punto de vista del estado mutable. Dos acciones en la misma rutina no pueden ser concurrentes. Y al igual que con los subprocesos, debe evitar compartir su estado mutable entre rutinas o tendrá que preocuparse por la sincronización. Evite compartir el estado mutable. Limite cada objeto mutable a un solo hilo o a una sola rutina y duerma bien.
Y esto es aplicable a los Flows , porque la recopilación de Flow funciona directamente en la corrutina que llama a un operador de terminal. No se lanzan nuevas rutinas por defecto.