Tengo un montón de columnas como matrices de cadenas de un archivo csv. Ahora quiero analizarlos. Dado que este análisis requiere un análisis de fechas y otras técnicas de análisis no tan rápidas, estaba pensando en el paralelismo (lo cronometré, lleva algo de tiempo). Mi enfoque simple:
Stream.of(columns).parallel().forEach(column -> result[column.index] = parseColumn(valueCache[column.index], column.type)); Las columnas contienen elementos ColumnDescriptor que simplemente tienen dos atributos, el índice de la columna que se analizará y el tipo que define cómo analizarlo. Nada más. resultado es una matriz de objetos que toma las matrices resultantes.
El problema ahora es que la función de análisis arroja una ParseException, que manejo más arriba en la pila de llamadas. Ya que estamos en paralelo aquí, no se puede simplemente lanzar. ¿Cuál es la mejor manera de manejar esto?
Tengo esta solución, pero me estremezco al leerla. ¿Cuál sería una mejor manera de hacerlo?
final CompletableFuture<ParseException> thrownException = new CompletableFuture<>(); Stream.of(columns).parallel().forEach(column -> { try { result[column.index] = parseColumn(valueCache[column.index], column.type); } catch (ParseException e) { thrownException.complete(e); }}); if(thrownException.isDone()) //only can be done if there is a value set. throw thrownException.getNow(null);Notas: No necesito todas las excepciones. Si los analizo secuencialmente, solo obtendré uno de todos modos. Así que está bien.
El problema es su premisa equivocada "Dado que estamos en paralelo aquí, no se puede descartar". No hay ninguna especificación que prohíba lanzar excepciones en el procesamiento paralelo. Simplemente puede lanzar esa excepción en una secuencia paralela de la misma manera que lo hace en una secuencia secuencial, envolviéndola en una excepción no verificada, si es una excepción verificada.
Si hay al menos una excepción lanzada en un hilo, la invocación forEach la propagará (o una de ellas) a la persona que llama.
El único problema que puede encontrar es que la implementación actual no espera a que se completen todos los subprocesos cuando encuentra una excepción. Esto se puede solucionar usando
try { Arrays.stream(columns).parallel() .forEach(column -> result[column.index] = parseColumn(valueCache[column.index], column.type)); } catch(Throwable t) { ForkJoinPool.commonPool().awaitQuiescence(1, TimeUnit.MINUTES); throw t; }Pero, por lo general, no lo necesita, ya que no accederá al resultado procesado simultáneamente en el caso excepcional.
Creo que la pregunta es más, ¿qué haces normalmente cuando lo analizas en serie?
¿Se detiene en la primera excepción y detiene todo el proceso? En ese caso, envuelva la excepción en una excepción de tiempo de ejecución y deje que la secuencia se cancele y la lance. Detecte la excepción del contenedor, desenvuélvala y manéjela.
¿Te saltas los malos registros? Luego, 1. realice un seguimiento de los errores en una Lista en algún lugar o 2. cree un objeto contenedor que pueda contener un resultado analizado o un error (no realice un seguimiento de las excepciones en sí, solo el mínimo necesario para describir el error).
Verifique luego si hubo errores en la lista para la primera opción, o muestre los registros que tuvieron errores de manera diferente para la segunda opción.