Play Framework 2.5文件处理进度的客户端通知方案问询
刚好我有过类似的Play Framework项目经验,来给你详细说说怎么实现实时进度更新——不管是你想的Akka BroadcastHub WebSocket方案,还是更轻量的替代方案,都给你安排明白!
方案一:基于Akka BroadcastHub的WebSocket(Java版)
这个方案适合需要双向通信的场景,我们可以调整成每个用户独享进度通道,避免所有人看到相同进度。
步骤1:定义进度消息模型
首先创建一个简单的POJO来封装进度数据:
public class ProgressEvent { public final int percentage; public final String message; public ProgressEvent(int percentage, String message) { this.percentage = percentage; this.message = message; } }
步骤2:实现WebSocket端点与上传处理
在你的Controller里,用Akka Stream创建每个上传任务专属的广播源,并关联到WebSocket:
import akka.actor.ActorSystem; import akka.stream.Materializer; import akka.stream.javadsl.BroadcastHub; import akka.stream.javadsl.Keep; import akka.stream.javadsl.Sink; import akka.stream.javadsl.Source; import play.mvc.Controller; import play.mvc.Result; import play.mvc.WebSocket; import javax.inject.Inject; import java.util.Map; import java.util.UUID; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionStage; import java.util.concurrent.ConcurrentHashMap; public class UploadController extends Controller { private final ActorSystem system; private final Materializer mat; private final Map<String, Source<ProgressEvent, ?>> taskProgressSources = new ConcurrentHashMap<>(); @Inject public UploadController(ActorSystem system, Materializer mat) { this.system = system; this.mat = mat; } // 带任务ID的WebSocket端点,只返回当前任务的进度 public WebSocket ws(String taskId) { return WebSocket.Json.accept(request -> { Source<ProgressEvent, ?> taskSource = taskProgressSources.get(taskId); if (taskSource == null) { // 返回无效任务ID的错误消息 return Source.single(new ProgressEvent(-1, "无效的任务ID,请重新上传")) .map(play.libs.Json::toJson) .mapMaterializedValue(source -> CompletableFuture.completedFuture(null)); } return taskSource.map(play.libs.Json::toJson) .mapMaterializedValue(source -> { // 客户端断开或任务完成后清理资源 source.whenComplete((done, err) -> taskProgressSources.remove(taskId)); return CompletableFuture.completedFuture(null); }); }); } // 文件上传接口,返回任务ID给前端 public CompletionStage<Result> upload() { return request().body().asMultipartFormData().map(multipart -> { return multipart.getFile("file").map(filePart -> { // 生成唯一任务ID String taskId = UUID.randomUUID().toString(); // 创建该任务专属的广播源和接收Sink Source<ProgressEvent, Sink<ProgressEvent, ?>> broadcastPair = Source.<ProgressEvent>actorRef(100, akka.stream.OverflowStrategy.dropNew()) .toMat(BroadcastHub.of(ProgressEvent.class, 256), Keep.both()) .run(mat); taskProgressSources.put(taskId, broadcastPair.first()); Sink<ProgressEvent, ?> progressSink = broadcastPair.second(); // 启动异步任务处理Excel CompletableFuture.runAsync(() -> { try { // 用Apache POI获取总行数(记得排除表头哦) int totalRows = getTotalExcelRows(filePart.getFile()); int processedRows = 0; int lastReportedPercent = -1; // 逐行处理Excel数据 processExcelRows(filePart.getFile(), row -> { processedRows++; int currentPercent = (int) Math.round((processedRows * 100.0) / totalRows); // 每10%更新一次进度(避免频繁推送) if (currentPercent % 10 == 0 && currentPercent != lastReportedPercent) { lastReportedPercent = currentPercent; // 发送进度消息到Sink progressSink.runWith(Source.single( new ProgressEvent(currentPercent, "已处理 " + processedRows + "/" + totalRows + " 行") ), mat); } // 执行你的自定义业务逻辑 handleRowData(row); }); // 处理完成后发送100%进度 progressSink.runWith(Source.single( new ProgressEvent(100, "文件处理完成!") ), mat); } catch (Exception e) { // 发送错误消息 progressSink.runWith(Source.single( new ProgressEvent(-1, "处理失败:" + e.getMessage()) ), mat); } finally { // 清理任务资源 taskProgressSources.remove(taskId); } }); // 返回任务ID给前端,用于连接WebSocket return ok(play.libs.Json.toJson(taskId)); }).orElseGet(() -> badRequest("请选择要上传的Excel文件")); }).orElseGet(() -> CompletableFuture.completedFuture(badRequest("请求格式错误"))); } // ---------------------- 以下是需要你实现的辅助方法 ---------------------- private int getTotalExcelRows(java.io.File excelFile) throws Exception { // 用Apache POI读取Excel总行数(示例:XSSFWorkbook) try (org.apache.poi.xssf.usermodel.XSSFWorkbook workbook = new org.apache.poi.xssf.usermodel.XSSFWorkbook(excelFile)) { return workbook.getSheetAt(0).getPhysicalNumberOfRows(); } } private void processExcelRows(java.io.File excelFile, java.util.function.Consumer<Object> rowHandler) throws Exception { // 用Apache POI逐行读取Excel,调用rowHandler处理每一行 try (org.apache.poi.xssf.usermodel.XSSFWorkbook workbook = new org.apache.poi.xssf.usermodel.XSSFWorkbook(excelFile)) { org.apache.poi.xssf.usermodel.XSSFSheet sheet = workbook.getSheetAt(0); for (org.apache.poi.ss.usermodel.Row row : sheet) { rowHandler.accept(row); } } } private void handleRowData(Object row) { // 你的自定义业务逻辑:比如校验数据、更新数据库等 } }
步骤3:前端实现WebSocket连接与进度展示
// 文件上传表单提交逻辑 document.getElementById("upload-form").addEventListener("submit", async function(e) { e.preventDefault(); const formData = new FormData(this); const uploadStatus = document.getElementById("upload-status"); const progressBar = document.getElementById("progress-bar"); const progressMsg = document.getElementById("progress-message"); try { const response = await fetch("/upload", { method: "POST", body: formData }); const taskId = await response.json(); uploadStatus.textContent = "文件上传成功,正在处理中..."; progressBar.style.width = "0%"; progressMsg.textContent = "初始化处理..."; // 连接对应任务的WebSocket const socket = new WebSocket(`ws://${window.location.host}/ws/${taskId}`); socket.onmessage = function(event) { const progress = JSON.parse(event.data); if (progress.percentage === -1) { // 处理错误 progressMsg.textContent = progress.message; socket.close(); } else { // 更新进度条和消息 progressBar.style.width = `${progress.percentage}%`; progressMsg.textContent = progress.message; if (progress.percentage === 100) { socket.close(); uploadStatus.textContent = "处理完成!"; } } }; socket.onerror = function(error) { console.error("WebSocket连接出错:", error); progressMsg.textContent = "连接进度服务器失败,请刷新页面重试"; }; socket.onclose = function() { console.log("进度连接已关闭"); }; } catch (err) { uploadStatus.textContent = "上传失败:" + err.message; } });
方案二:轻量替代——Server-Sent Events(SSE)
如果不需要双向通信,只是服务器单向推送进度,SSE会更简单,不需要处理WebSocket的双向逻辑,代码量更少。
步骤1:Controller实现SSE端点与上传处理
import play.mvc.Controller; import play.mvc.Result; import play.libs.EventSource; import akka.stream.javadsl.Source; import javax.inject.Inject; import java.util.Map; import java.util.UUID; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionStage; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentLinkedQueue; public class UploadController extends Controller { private final Map<String, ConcurrentLinkedQueue<String>> taskProgressQueues = new ConcurrentHashMap<>(); @Inject public UploadController() {} // SSE端点,根据任务ID返回进度 public Result sse(String taskId) { ConcurrentLinkedQueue<String> queue = taskProgressQueues.get(taskId); if (queue == null) { return ok(EventSource.event("{\"percentage\": -1, \"message\": \"无效的任务ID\"}")).as("text/event-stream"); } Source<String, ?> source = Source.fromIterator(() -> queue.iterator()) .concat(Source.maybe()) // 保持连接直到任务完成 .map(EventSource::event); return ok().chunked(source).as("text/event-stream"); } // 文件上传接口 public CompletionStage<Result> upload() { return request().body().asMultipartFormData().map(multipart -> { return multipart.getFile("file").map(filePart -> { String taskId = UUID.randomUUID().toString(); ConcurrentLinkedQueue<String> progressQueue = new ConcurrentLinkedQueue<>(); taskProgressQueues.put(taskId, progressQueue); CompletableFuture.runAsync(() -> { try { int totalRows = getTotalExcelRows(filePart.getFile()); int processedRows = 0; int lastReportedPercent = -1; processExcelRows(filePart.getFile(), row -> { processedRows++; int currentPercent = (int) Math.round((processedRows * 100.0) / totalRows); if (currentPercent % 10 == 0 && currentPercent != lastReportedPercent) { lastReportedPercent = currentPercent; progressQueue.add(String.format( "{\"percentage\": %d, \"message\": \"已处理 %d/%d 行\"}", currentPercent, processedRows, totalRows )); } handleRowData(row); }); progressQueue.add("{\"percentage\": 100, \"message\": \"文件处理完成!\"}"); } catch (Exception e) { progressQueue.add(String.format( "{\"percentage\": -1, \"message\": \"处理失败:%s\"}", e.getMessage() )); } finally { // 延迟清理,确保客户端收到最后一条消息 CompletableFuture.delayedExecutor(5, java.util.concurrent.TimeUnit.SECONDS).execute(() -> { taskProgressQueues.remove(taskId); }); } }); return ok(play.libs.Json.toJson(taskId)); }).orElseGet(() -> badRequest("请选择要上传的Excel文件")); }).orElseGet(() -> CompletableFuture.completedFuture(badRequest("请求格式错误"))); } // 辅助方法同方案一,这里省略... private int getTotalExcelRows(java.io.File excelFile) throws Exception { /* ... */ } private void processExcelRows(java.io.File excelFile, java.util.function.Consumer<Object> rowHandler) throws Exception { /* ... */ } private void handleRowData(Object row) { /* ... */ } }
步骤2:前端SSE实现
document.getElementById("upload-form").addEventListener("submit", async function(e) { e.preventDefault(); const formData = new FormData(this); const uploadStatus = document.getElementById("upload-status"); const progressBar = document.getElementById("progress-bar"); const progressMsg = document.getElementById("progress-message"); try { const response = await fetch("/upload", { method: "POST", body: formData }); const taskId = await response.json(); uploadStatus.textContent = "文件上传成功,正在处理中..."; progressBar.style.width = "0%"; progressMsg.textContent = "初始化处理..."; // 连接SSE const eventSource = new EventSource(`/sse/${taskId}`); eventSource.onmessage = function(event) { const progress = JSON.parse(event.data); if (progress.percentage === -1) { progressMsg.textContent = progress.message; eventSource.close(); } else { progressBar.style.width = `${progress.percentage}%`; progressMsg.textContent = progress.message; if (progress.percentage === 100) { eventSource.close(); uploadStatus.textContent = "处理完成!"; } } }; eventSource.onerror = function(error) { console.error("SSE连接出错:", error); progressMsg.textContent = "连接进度服务器失败,请刷新页面重试"; eventSource.close(); }; } catch (err) { uploadStatus.textContent = "上传失败:" + err.message; } });
关键注意事项
- 进度计算准确性:记得排除Excel的表头行,避免总行数计算错误导致进度显示偏差。
- 资源清理:不管是WebSocket还是SSE,任务完成或客户端断开后一定要清理对应的资源(比如从Map中移除任务ID),避免内存泄漏。
- 异常处理:Excel处理过程中要捕获所有可能的异常,及时给用户返回错误消息,避免任务静默失败。
- 并发安全:用
ConcurrentHashMap和ConcurrentLinkedQueue来存储任务进度,确保多线程环境下的安全性。
内容的提问来源于stack exchange,提问作者Omnia89
相关产品推荐
相关产品推荐

