Tengo algunos AsyncEnumerable<string> s que me gustaría fusionar en un solo AsyncEnumerable<string> , que debe contener todos los elementos que se emiten simultáneamente desde esas secuencias. Así que utilicé el operador Merge del paquete System.Interactive.Async . El problema es que este operador no siempre trata todas las secuencias como iguales. En algunas circunstancias, prefiere emitir elementos de las secuencias que están en el lado izquierdo de la lista de argumentos y descuida las secuencias que están en el lado derecho de la lista de argumentos. Aquí hay un ejemplo mínimo que reproduce este comportamiento indeseable:
var sequence_A = Enumerable.Range(1, 5).Select(i => $"A{i}").ToAsyncEnumerable(); var sequence_B = Enumerable.Range(1, 5).Select(i => $"B{i}").ToAsyncEnumerable(); var sequence_C = Enumerable.Range(1, 5).Select(i => $"C{i}").ToAsyncEnumerable(); var merged = AsyncEnumerableEx.Merge(sequence_A, sequence_B, sequence_C); await foreach (var item in merged) Console.WriteLine(item); Este fragmento de código también depende del paquete System.Linq.Async . La sequence_A emite 5 elementos a partir de "A" , la sequence_B emite 5 elementos a partir de "B" y la sequence_C emite 5 elementos a partir de "C" .
Salida (no deseable):
A1 A2 A3 A4 A5 B1 B2 B3 B4 B5 C1 C2 C3 C4 C5La salida deseable debería verse así:
A1 B1 C1 A2 B2 C2 A3 B3 C3 A4 B4 C4 A5 B5 C5En caso de que todas las secuencias tengan su siguiente elemento disponible, la secuencia fusionada debe extraer un elemento de cada secuencia, en lugar de extraer elementos repetidamente de la secuencia más a la izquierda.
¿Cómo puedo asegurarme de que mis secuencias se fusionen con equidad? Estoy buscando una combinación de operadores de los paquetes oficiales que tenga el comportamiento deseable, o un operador Merge personalizado que haga lo que quiero.
Aclaración: estoy interesado en la funcionalidad Merge concurrente , donde todas las secuencias de origen se observan al mismo tiempo y cualquier emisión de cualquiera de las secuencias se propaga a la secuencia fusionada. El concepto de equidad se aplica cuando más de una secuencia puede emitir un elemento inmediatamente, en cuyo caso sus emisiones deben estar intercaladas. En el caso contrario, cuando no hay ningún elemento inmediatamente disponible, la regla es "primero en llegar, primero en irse".
Actualización: aquí hay una demostración más realista, que incluye la latencia en las secuencias del productor y en el ciclo de enumeración de consumo. Simula una situación en la que consumir los valores producidos por la secuencia más a la izquierda toma más tiempo que el tiempo requerido para producir esos valores.
var sequence_A = Produce("A", 200, 1, 2, 3, 4, 5); var sequence_B = Produce("B", 150, 1, 2, 3, 4, 5); var sequence_C = Produce("C", 100, 1, 2, 3, 4, 5); var merged = AsyncEnumerableEx.Merge(sequence_A, sequence_B, sequence_C); await foreach (var item in merged) { Console.WriteLine(item); await Task.Delay(item.StartsWith("A") ? 300 : 50); // Latency } async IAsyncEnumerable<string> Produce(string prefix, int delay, params int[] values) { foreach (var value in values) { var delayTask = Task.Delay(delay); yield return $"{prefix}{value}"; await delayTask; // Latency } } El resultado es un sesgo no deseado para los valores producidos por la sequence_A :
A1 A2 A3 A4 A5 B1 B2 C1 B3 C2 B4 C3 C4 B5 C5El ejemplo es un poco artificial ya que todos los resultados están disponibles de inmediato. Si se agrega incluso un pequeño retraso, los resultados son mixtos:
var sequence_A = AsyncEnumerable.Range(1, 5) .SelectAwait(async i =>{ await Task.Delay(i); return $"A{i}";}); var sequence_B = AsyncEnumerable.Range(1, 5) .SelectAwait(async i =>{ await Task.Delay(i); return $"B{i}";}); var sequence_C = AsyncEnumerable.Range(1, 5) .SelectAwait(async i =>{ await Task.Delay(i); return $"C{i}";}); var sequence_D = AsyncEnumerable.Range(1, 5) .SelectAwait(async i =>{ await Task.Delay(i); return $"D{i}";}); await foreach (var item in seq) Console.WriteLine(item);Esto produce resultados diferentes y mixtos cada vez:
B1 A1 C1 D1 D2 A2 B2 C2 D3 A3 B3 C3 C4 A4 B4 D4 D5 A5 B5 C5Los comentarios del método explican que se volvió a implementar para que fuera más barato y más justo:
// // This new implementation of Merge differs from the original one in a few ways: // // - It's cheaper because: // - no conversion from ValueTask<bool> to Task<bool> takes place using AsTask, // - we don't instantiate Task.WhenAny tasks for each iteration. // - It's fairer because: // - the MoveNextAsync tasks are awaited concurently, but completions are queued, // instead of awaiting a new WhenAny task where "left" sources have preferential // treatment over "right" sources. //