Tengo un ciclo foreach que es responsable de ejecutar un determinado conjunto de declaraciones. Una parte de eso es guardar una imagen de una URL en el almacenamiento de Azure. Tengo que hacer esto para un gran conjunto de datos. Para lograr lo mismo, he convertido el bucle foreach en un bucle Parallel.ForEach .
Parallel.ForEach(listSkills, item => { // some business logic var b = getImageFromUrl(item.Url); Stream ms = new MemoryStream(b); saveImage(ms); // more business logic }); private static byte[] getByteArray(Stream input) { using (MemoryStream ms = new MemoryStream()) { input.CopyTo(ms); return ms.ToArray(); } } public static byte[] getImageFromUrl(string url) { HttpWebRequest request = null; HttpWebResponse response = null; byte[] b = null; request = (HttpWebRequest)WebRequest.Create(url); response = (HttpWebResponse)request.GetResponse(); if (request.HaveResponse) { if (response.StatusCode == HttpStatusCode.OK) { Stream receiveStream = response.GetResponseStream(); b = getByteArray(receiveStream); } } return b; } public static void saveImage(Stream fileContent) { fileContent.Seek(0, SeekOrigin.Begin); byte[] bytes = getByteArray(fileContent); var blob = null; blob.UploadFromByteArrayAsync(bytes, 0, bytes.Length).Wait(); }Aunque hay casos en los que recibo el siguiente error y la imagen no se guarda.
El host remoto cerró a la fuerza una conexión existente.
También compartiendo el StackTrace:
at System.Net.Sockets.NetworkStream.Read(Span`1 buffer) at System.Net.Security.SslStream.<FillBufferAsync>d__183`1.MoveNext() at System.Runtime.ExceptionServices.ExceptionDispatchInfo.Throw() at System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task) at System.Net.Security.SslStream.<ReadAsyncInternal>d__181`1.MoveNext() at System.Runtime.ExceptionServices.ExceptionDispatchInfo.Throw() at System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task) at System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task) at System.Net.Security.SslStream.Read(Byte[] buffer, Int32 offset, Int32 count) at System.IO.Stream.Read(Span`1 buffer) at System.Net.Http.HttpConnection.Read(Span`1 destination) at System.Net.Http.HttpConnection.ContentLengthReadStream.Read(Span`1 buffer) at System.Net.Http.HttpBaseStream.Read(Byte[] buffer, Int32 offset, Int32 count) at System.IO.Stream.CopyTo(Stream destination, Int32 bufferSize) at Utilities.getByteArray(Stream input) in D:\repos\SampleProj\Sample\Helpers\CH.cs:line 238 at Utilities.getImageFromUrl(String url) in D:\repos\SampleProj\Sample\Helpers\CH.cs:line 178 Supongo que esto tal vez porque no estoy usando cerraduras. No estoy seguro de si usar bloqueos dentro de un bucle Parallel.ForEach .
De acuerdo con otra pregunta sobre stackoverflow , estas son las causas potenciales de que el host remoto cerró a la fuerza una conexión existente. :
- Está enviando datos con formato incorrecto a la aplicación (lo que podría incluir el envío de una solicitud HTTPS a un servidor HTTP)
- El enlace de red entre el cliente y el servidor está fallando por alguna razón
- Ha desencadenado un error en la aplicación de terceros que provocó que se bloqueara
- La aplicación de terceros ha agotado los recursos del sistema
Dado que solo algunas de sus solicitudes se ven afectadas, creo que podemos excluir la primera. Esto puede ser, por supuesto, un problema de red y, en ese caso, esto sucederá de vez en cuando dependiendo de la calidad de la red entre usted y el servidor.
A menos que encuentre una indicación de un error de AzureStorage de otros usuarios, existe una alta probabilidad de que su llamada esté consumiendo demasiados recursos del servidor remoto (conexiones/datos) al mismo tiempo. Los servidores y el proxy tienen limitaciones sobre la cantidad de conexiones que pueden manejar al mismo tiempo (especialmente desde la misma máquina cliente).
Dependiendo del tamaño de su lista listSkills , su código puede lanzar una gran cantidad de solicitudes en paralelo (tantas como pueda hacerlo su grupo de subprocesos), posiblemente inundando el servidor.
Al menos podría limitar el número de lanzamientos de tareas paralelas utilizando MaxDegreeOfParallelism de esta manera:
Parallel.ForEach(listSkills, new ParallelOptions { MaxDegreeOfParallelism = 4 }, item => { // some business logic var b = getImageFromUrl(item.Url); Stream ms = new MemoryStream(b); saveImage(ms); // more business logic });Puede controlar el paralelismo como:
listSkills.AsParallel() .Select(item => {/*Your Logic*/ return item}) .WithDegreeOfParallelism(10) .Select(item => { getImageFromUrl(item.url); saveImage(your_stream); return item; }); Pero Parallel.ForEach no es bueno para IO porque está diseñado para tareas CPU-intensive , si lo usa para operaciones IO-bound especialmente para realizar solicitudes web, puede desperdiciar hilos de grupo de subprocesos bloqueados mientras espera una respuesta.
Utiliza métodos de solicitud web asincrónicos como HttpWebRequest.GetResponseAsync , por otro lado, también puede usar construcciones de sincronización de subprocesos para eso, como un ejemplo ab usando Semaphore , el Semaphore es como una cola, permite que pasen X subprocesos, y el resto debe esperar hasta uno de los hilos ocupados terminará su trabajo. Primero haga que su método getStream sea async (esta no es una buena solución, pero puede ser mejor):
public static async Task getImageFromUrl(SemaphoreSlim semaphore, string url) { try { HttpWebRequest request = null; byte[] b = null; request = (HttpWebRequest)WebRequest.Create(url); using (var response = await request.GetResponseAsync().ConfigureAwait(false)) { // your logic } } catch (Exception ex) { // handle exp } finally { // release semaphore.Release(); } }y luego:
using (var semaphore = new SemaphoreSlim(10)) { foreach (var url in urls) { // await here until there is a room for this task await semaphore.WaitAsync(); tasks.Add(getImageFromUrl(semaphore, url)); } // await for the rest of tasks to complete await Task.WhenAll(tasks); } No debe usar Parallel o Task.Run en su lugar, puede tener un método de controlador async como:
public async Task handleResponse(Task<HttpResponseMessage> response) { HttpResponseMessage response = await response; //Process your data } y luego use Task.WhenAll como:
Task[] requests = myList.Select(l => getImageFromUrl(l.Id)) .Select(r => handleResponse(r)) .ToArray(); await Task.WhenAll(requests); al final, hay varias soluciones para su escenario, pero olvídese de Parallel.Foreach . Foreach, en su lugar, use una solución optimizada.
Hay varios problemas con este código:
Parallel.ForEach está diseñado para el paralelismo de datos, no para IO. El código está congelando todos los núcleos de la CPU esperando que IO se completeHay varias formas de ejecutar muchas operaciones de IO al mismo tiempo en .NET Core.
.NET 6
En la versión actual de soporte a largo plazo de .NET, .NET 6, esto se puede hacer usando Parallel.ForEachAsync . Scott Hanselman muestra lo fácil que es usarlo para llamadas API
Puede recuperar los datos directamente con GetBytesAsync :
record CopyRequest(Uri sourceUri,Uri blobUri); ... var requests=new List<CopyRequest>(); //Load some source/target URLs var client=new HttpClient(); await Parallel.ForEachAsync(requests,async req=>{ var bytes=await client.GetBytesAsync(req.sourceUri); var blob=new CloudAppendBlob(req.targetUri); await blob.UploadFromByteArrayAsync(bytes, 0, bytes.Length); });Una mejor opción sería recuperar los datos como una secuencia y enviarlos directamente al blob:
await Parallel.ForEachAsync(requests,async req=>{ var response=await client.GetAsync(req.sourceUri, HttpCompletionOption.ResponseHeadersRead); using var sourceStream=await response.Content.ReadAsStreamAsync(); var blob=new CloudAppendBlob(req.targetUri); await blob.UploadFromStreamAsync(sourceStream); }); HttpCompletionOption.ResponseHeadersRead hace que GetAsync regrese tan pronto como se reciban los encabezados de respuesta, sin almacenar en búfer ninguno de los datos de respuesta.
.NET 3.1
En versiones anteriores de .NET Core (que llegarán al final de su vida útil en unos meses), puede usar, por ejemplo, un ActionBlock con un grado de paralelismo superior a 1:
var options=new ExecuteDataflowBlockOptions{ MaxDegreeOfParallelism = 8}; var copyBlock=new ActionBlock<CopyRequest>(async req=>{ var response=await client.GetAsync(req.sourceUri, HttpCompletionOption.ResponseHeadersRead); using var sourceStream=await response.Content.ReadAsStreamAsync(); var blob=new CloudAppendBlob(req.targetUri); await blob.UploadFromStreamAsync(sourceStream); }, options);Las clases de bloque en la biblioteca TPL Dataflow se pueden usar para construir canalizaciones de procesamiento similares a una canalización de script de shell, con cada bloque canalizando su salida al siguiente bloque.