You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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;
    }
});

关键注意事项
  1. 进度计算准确性:记得排除Excel的表头行,避免总行数计算错误导致进度显示偏差。
  2. 资源清理:不管是WebSocket还是SSE,任务完成或客户端断开后一定要清理对应的资源(比如从Map中移除任务ID),避免内存泄漏。
  3. 异常处理:Excel处理过程中要捕获所有可能的异常,及时给用户返回错误消息,避免任务静默失败。
  4. 并发安全:用ConcurrentHashMap和ConcurrentLinkedQueue来存储任务进度,确保多线程环境下的安全性。

内容的提问来源于stack exchange,提问作者Omnia89

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.29 07:54:11