MicroProfile EE8应用Worker线程代码未执行,实际跑在原线程问题排查
我在Payara Server上开发Java EE 8 MicroProfile应用,需要对REST API接收的XLSX文件做并行处理。使用ManagedExecutorService(也尝试过newFixedThreadPool)创建Callable执行并行任务,但出现异常行为:Callable确实在Worker线程中运行,但其内部调用的ApplicationScoped Bean RowConverter的业务逻辑却跑在HTTP线程池里。所有相关Bean均为ApplicationScoped,已排除Singleton EJB的影响,问题依然存在。
@ApplicationScoped public class BusinessService { private @Inject RowConverter rowConverter; private @Resource ManagedExecutorService managedExecutorService; public Map<RequestValidity, List<EInvoiceRequest>> process(InputStream fileInputStream) { try (Workbook workbook = WorkbookFactory.create(fileInputStream)) { Sheet sheet = workbook.getSheetAt(0); Iterator<Row> iterator = IteratorUtils.filteredIterator(sheet.rowIterator(), row -> row.getRowNum() > 0 && row.getCell("FOOBAR", Row.MissingCellPolicy.RETURN_BLANK_AS_NULL) != null); List<Row> rows = IteratorUtils.toList(iterator); List<Callable<Stream<MyWrapperObj>>> callables = ListUtils.partition(rows, Math.max(300, rows.size() / 10)).stream() .map(partition -> new RowToMyWrapperObjCallable(partition, rowConverter)) .collect(Collectors.toList()); return managedExecutorService.invokeAll(callables).stream() .flatMap(future -> { try { return future.get(); } catch (InterruptedException e) { log.error("Error processing XLSX file", e); Thread.currentThread().interrupt(); } catch (ExecutionException e) { log.error("Error processing XLSX file", e); } throw new GenericException("Error processing commissions XLSX file", ErrorCode.FILE_PROCESSING_ERROR); }).collect(Collectors.groupingBy(r -> isRequestValid(r) ? RequestValidity.VALID : RequestValidity.INVALID)); } catch (Exception e) { throw new GenericException("Error processing XLSX file", ErrorCode.FILE_PROCESSING_ERROR); } } } // ------------------------------------------------------------------------------------------------ // @RequiredArgsConstructor public class RowToMyWrapperObjCallable implements Callable<Stream<MyWrapperObj>> { private final List<Row> rows; private final RowConverter mapper; @Override public Stream<MyWrapperObj> call() throws Exception { log.info("Running mapping on thread: {}", Thread.currentThread().getName()); return rows.stream().map(mapper); } } // ------------------------------------------------------------------------------------------------ // @ApplicationScoped public class RowConverter implements Function<Row, MyWrapperObj> { @Override public MyWrapperObj apply(Row row) { log.info("[{}] RowConverter::apply", Thread.currentThread().getName()); // 执行映射逻辑 return new MyWrapperObj(); } }
2024-12-06 00:46:57.968+0200 INFO RowToMyWrapperObjCallable.call(Line 28) [pool-59-thread-1](523) [] Running mapping on thread: pool-59-thread-1 2024-12-06 00:46:57.968+0200 INFO RowToMyWrapperObjCallable.call(Line 28) [pool-59-thread-3](525) [] Running mapping on thread: pool-59-thread-3 2024-12-06 00:46:57.969+0200 INFO RowToMyWrapperObjCallable.call(Line 28) [pool-59-thread-4](526) [] Running mapping on thread: pool-59-thread-4 2024-12-06 00:46:57.968+0200 INFO RowToMyWrapperObjCallable.call(Line 28) [pool-59-thread-2](524) [] Running mapping on thread: pool-59-thread-2 2024-12-06 00:47:00.633+0200 INFO RowConverter.apply(Line 60) [http-thread-pool::http-listener-1(1)](149) [] [http-thread-pool::http-listener-1(1)] RowConverter::apply 2024-12-06 00:47:00.889+0200 INFO RowConverter.apply(Line 60) [http-thread-pool::http-listener-1(1)](149) [] [http-thread-pool::http-listener-1(1)] RowConverter::apply 2024-12-06 00:47:01.084+0200 INFO RowConverter.apply(Line 60) [http-thread-pool::http-listener-1(1)](149) [] [http-thread-pool::http-listener-1(1)] RowConverter::apply
核心原因是Stream的延迟执行特性:RowToMyWrapperObjCallable.call()方法返回的Stream<MyWrapperObj>只是一个"计算蓝图",其中的map(mapper)属于中间操作,不会立即执行实际的映射逻辑。只有当调用终端操作(比如collect、forEach)时,Stream才会触发计算。
在你的代码中,invokeAll获取Future后,future.get()拿到的只是未执行的Stream,后续的flatMap和collect操作是在HTTP请求线程中执行的,这才会真正触发RowConverter.apply(),所以日志显示该方法运行在HTTP线程池。
要确保并行处理生效,必须让映射逻辑在Worker线程中执行,核心是在Callable内部完成Stream的终端操作:
方案1:在Callable内完成映射并返回集合
修改RowToMyWrapperObjCallable的实现,提前执行终端操作将Stream转为List,确保映射逻辑在Worker线程执行:
@RequiredArgsConstructor public class RowToMyWrapperObjCallable implements Callable<List<MyWrapperObj>> { private final List<Row> rows; private final RowConverter mapper; @Override public List<MyWrapperObj> call() throws Exception { log.info("Running mapping on thread: {}", Thread.currentThread().getName()); return rows.stream().map(mapper).collect(Collectors.toList()); } }
同时调整BusinessService中的泛型和后续处理逻辑:
// 调整Callable泛型为List<MyWrapperObj> List<Callable<List<MyWrapperObj>>> callables = ListUtils.partition(rows, Math.max(300, rows.size() / 10)).stream() .map(partition -> new RowToMyWrapperObjCallable(partition, rowConverter)) .collect(Collectors.toList()); return managedExecutorService.invokeAll(callables).stream() .flatMap(future -> { try { // 将List转为Stream继续后续分组 return future.get().stream(); } catch (InterruptedException e) { log.error("Error processing XLSX file", e); Thread.currentThread().interrupt(); } catch (ExecutionException e) { log.error("Error processing XLSX file", e); } throw new GenericException("Error processing commissions XLSX file", ErrorCode.FILE_PROCESSING_ERROR); }).collect(Collectors.groupingBy(r -> isRequestValid(r) ? RequestValidity.VALID : RequestValidity.INVALID));
方案2:使用并行Stream(可选)
如果希望利用Stream的并行能力,可在Callable内直接使用parallelStream,但需注意Java EE环境中,并行Stream默认使用ForkJoinPool,若要绑定到ManagedExecutorService,需要额外配置上下文传递(Payara支持@ContextService绑定),方案1更简单直接。
内容的提问来源于stack exchange,提问作者George Karanikas

