找回密码
立即注册
搜索
热搜: Java Python Linux Go
发回帖 发新帖
Claude、GPT 海外模型 API 接入云原生前端项目实战教程50G互联网架构师面试指南
大模型全栈开发课程企业级DevOps全栈实践零基础产品经理就业课程

4685

积分

0

好友

597

主题
发表于 1 小时前 | 查看: 3| 回复: 0

前段时间改了一个后台的 Excel 导入功能。

最开始这个功能其实很普通。

运营上传一个 Excel,里面是商品数据,后端解析以后写进 MySQL

数据少的时候完全没问题。

几百条:

1 秒左右

几千条:

几秒

后来运营拿了一份 10 万多行的 Excel 上来。

页面点完“导入”以后,一直转圈。

转了几十秒。

最后浏览器直接提示:

504 Gateway Timeout

运营又点了一次。

结果更麻烦。

第一批任务其实还在后台跑,第二批数据又开始导。

数据库里出现了一堆重复数据。

我去看代码,发现原来的实现大概是这样的:

@PostMapping(”/import”)
public ImportResult importExcel(
        @RequestParam(”file”) MultipartFile file)
        throws Exception {

    List rows =
            excelReader.read(
                    file.getInputStream()
            );

    List products =
            rows.stream()
                    .map(this::convert)
                    .toList();

    productRepository.saveAll(products);

    return ImportResult.success(
            products.size()
    );
}

几千条的时候,这段代码一点毛病都看不出来。

但数据一到十万级,几个问题同时出来了。

Excel 全部读进内存。

10 万个 Java 对象留在一个 List 里。

HTTP 请求一直等数据库处理完成。

某一行数据错误,整批怎么处理也没想清楚。

用户根本不知道:

到底导到哪了?

更麻烦的是,只要网关超时,不代表后端任务真的停止。

用户看到的是:

导入失败

服务器实际上可能还在:

INSERT
INSERT
INSERT

这才是这类接口真正危险的地方。

后来我没有继续优化这个 Controller。

而是直接把整个模型换掉了。

现在上传 Excel 以后,接口不会负责完成导入。

它只做三件事:

保存上传文件
创建导入任务
返回 taskId

真正解析 Excel 和写数据库,在后台慢慢处理。

于是一个 Excel 导入从:

HTTP 请求
   ↓
解析 10 万行
   ↓
校验
   ↓
写数据库
   ↓
返回结果

变成了:

HTTP 请求
   ↓
保存文件
   ↓
创建任务
   ↓
立即返回 taskId

后台任务
   ↓
流式解析
   ↓
逐行校验
   ↓
500 条一批入库
   ↓
更新进度
   ↓
SSE 实时推送进度
   ↓
记录失败数据

这个改动做完以后,Controller 反而非常简单。

我先建了一张导入任务表。

CREATE TABLE import_task (
    id VARCHAR(64) PRIMARY KEY,
    file_name VARCHAR(255) NOT NULL,
    file_path VARCHAR(512) NOT NULL,
    status VARCHAR(32) NOT NULL,
    processed_rows BIGINT NOT NULL DEFAULT 0,
    success_rows BIGINT NOT NULL DEFAULT 0,
    failed_rows BIGINT NOT NULL DEFAULT 0,
    error_message VARCHAR(1000),
    created_at DATETIME NOT NULL,
    started_at DATETIME,
    finished_at DATETIME
);

状态我只用了几个:

WAITING
RUNNING
SUCCESS
PARTIAL_SUCCESS
FAILED

然后上传接口改成这样:

@PostMapping(”/api/imports”)
public ImportTaskResponse upload(
        @RequestParam(”file”)
        MultipartFile file)
        throws IOException {

    String taskId =
            UUID.randomUUID().toString();

    Path filePath =
            importFileStore.save(
                    taskId,
                    file
            );

    importTaskService.create(
            taskId,
            file.getOriginalFilename(),
            filePath.toString()
    );

    importTaskExecutor.execute(
            () -> importProcessor.process(
                    taskId,
                    filePath
            )
    );

    return new ImportTaskResponse(
            taskId,
            ”WAITING”
    );
}

这里有一个坑我以前也踩过。

千万不要这样:

importTaskExecutor.execute(
        () -> {
            file.getInputStream();
        }
);

也就是:

Controller 收到 MultipartFile
↓
直接把 MultipartFile 扔给异步线程

看起来没问题。

但 MultipartFile 本质上属于当前 HTTP 请求。

请求结束以后,Servlet 容器用于保存上传内容的临时资源可能已经被清理。

所以我现在第一步一定是:

Files.copy(...)

把文件保存到我们自己控制的目录。

例如:

/data/import/
    7c8290....xlsx

请求结束以后,这个文件还在。

后台任务什么时候处理都没有关系。

文件保存逻辑很简单:

@Component
public class ImportFileStore {

    private final Path root =
            Paths.get(”/data/import”);

    public Path save(
            String taskId,
            MultipartFile file)
            throws IOException {

        Files.createDirectories(root);

        Path target =
                root.resolve(
                        taskId + ”.xlsx”
                );

        try (InputStream input =
                     file.getInputStream()) {

            Files.copy(
                    input,
                    target,
                    StandardCopyOption
                            .REPLACE_EXISTING
            );
        }

        return target;
    }
}

异步线程池我也没有直接用默认的 @Async,而是专门给 Excel 导入建了一个线程池。

@Configuration
public class ImportExecutorConfig {

    @Bean
    public ThreadPoolTaskExecutor
            importTaskExecutor() {

        ThreadPoolTaskExecutor executor =
                new ThreadPoolTaskExecutor();

        executor.setCorePoolSize(2);
        executor.setMaxPoolSize(4);
        executor.setQueueCapacity(20);

        executor.setThreadNamePrefix(
                ”excel-import-”
        );

        executor.initialize();
        return executor;
    }
}

Excel 导入本身就是:

文件 IO
+
数据校验
+
数据库写入

开几十个任务同时导,并不会让一份 Excel 更快。

反而很容易把数据库连接池吃满。

所以这种后台批处理任务,我宁愿:

少量并发
+
排队

也不会无限开线程。

接下来就是 Excel 怎么读。

这里我这次没有继续照很多旧文章里的写法直接引 EasyExcel。

原因很现实:Alibaba EasyExcel 的 GitHub 仓库已经在 2025 年 9 月归档并进入只读状态,老项目如果已经稳定使用当然没必要马上替换,但新项目我更愿意选仍在持续维护的 Apache POI。

依赖:

org.apache.poi
    poi-ooxml
    5.5.1

以前不少代码是这样:

try (Workbook workbook =
             WorkbookFactory.create(file)) {

    Sheet sheet =
            workbook.getSheetAt(0);

    for (Row row : sheet) {
        ...
    }
}

小文件这么写没什么问题。

但大文件我现在会直接走 Event 模式。

核心区别是:

Workbook 模型:

把大量 Excel 对象构造成 Java 对象
然后遍历

SAX/Event:

读到一行
处理一行
然后丢掉

也就是说,10 万行并不意味着 JVM 里同时存在 10 万个 ProductImportRow。

我只保留当前这一批。

我最终做了一个流式解析器:

@Component
public class ExcelStreamReader {

    public void read(
            Path file,
            Consumer consumer)
            throws Exception {

        try (OPCPackage pkg =
                     OPCPackage.open(
                             file.toFile(),
                             PackageAccess.READ
                     )) {

            ReadOnlySharedStringsTable strings =
                    new ReadOnlySharedStringsTable(pkg);
            XSSFReader reader =
                    new XSSFReader(pkg);
            StylesTable styles =
                    reader.getStylesTable();
            DataFormatter formatter =
                    new DataFormatter();
            Iterator sheets =
                    reader.getSheetsData();
            if (!sheets.hasNext()) {
                return;
            }
            try (InputStream sheet =
                         sheets.next()) {
                XMLReader parser =
                        XMLHelper.newXMLReader();
                ProductSheetHandler handler =
                        new ProductSheetHandler(
                                consumer
                        );
                ContentHandler contentHandler =
                        new XSSFSheetXMLHandler(
                                styles,
                                null,
                                strings,
                                handler,
                                formatter,
                                false
                        );
                parser.setContentHandler(
                        contentHandler
                );
                parser.parse(
                        new InputSource(sheet)
                );
            }
        }
    }
}

真正收到一行数据的是 ProductSheetHandler:

public class ProductSheetHandler
        implements XSSFSheetXMLHandler
        .SheetContentsHandler {

    private final Consumer
            consumer;
    private final Map cells =
            new HashMap<>();
    public ProductSheetHandler(
            Consumer consumer) {
        this.consumer = consumer;
    }
    @Override
    public void startRow(int rowNum) {
        cells.clear();
    }
    @Override
    public void cell(
            String cellReference,
            String formattedValue,
            XSSFComment comment) {
        int column =
                new CellReference(
                        cellReference
                ).getCol();
        cells.put(
                column,
                formattedValue
        );
    }
    @Override
    public void endRow(int rowNum) {
        if (rowNum == 0) {
            return;
        }
        ProductImportRow row =
                new ProductImportRow(
                        rowNum + 1,
                        cells.get(0),
                        cells.get(1),
                        cells.get(2),
                        cells.get(3)
                );
        consumer.accept(row);
    }
}

这里比较关键的不是 POI 这几个类怎么写。

真正重要的是:

consumer.accept(row);

每解析一行就交出去。

而不是:

List rows =
        new ArrayList<>();

rows.add(...);

最后让 10 万、20 万甚至 50 万个对象一直留在内存里。

真正的导入逻辑,我只保留 500 条。

@Component
public class ImportProcessor {

    private static final int BATCH_SIZE = 500;
    private final ExcelStreamReader excelReader;
    private final ProductBatchWriter batchWriter;
    private final ImportTaskService taskService;
    private final ImportErrorService errorService;
    private final ImportProgressPublisher publisher;

    public void process(
            String taskId,
            Path file) {
        taskService.markRunning(taskId);
        List batch =
                new ArrayList<>(BATCH_SIZE);
        AtomicLong processed =
                new AtomicLong();
        AtomicLong success =
                new AtomicLong();
        AtomicLong failed =
                new AtomicLong();
        try {
            excelReader.read(
                    file,
                    row -> {
                        long current =
                                processed
                                        .incrementAndGet();
                        try {
                            Product product =
                                    validateAndConvert(
                                            row
                                    );
                            batch.add(product);
                        } catch (Exception e) {
                            failed.incrementAndGet();
                            errorService.save(
                                    taskId,
                                    row.rowNo(),
                                    row,
                                    e.getMessage()
                            );
                        }
                        if (batch.size()
                                >= BATCH_SIZE) {
                            int count =
                                    batchWriter
                                            .write(batch);
                            success.addAndGet(
                                    count
                            );
                            batch.clear();
                        }
                        if (current % 500 == 0) {
                            taskService
                                    .updateProgress(
                                            taskId,
                                            processed.get(),
                                            success.get(),
                                            failed.get()
                                    );
                            publisher.publish(
                                    taskId,
                                    processed.get(),
                                    success.get(),
                                    failed.get()
                            );
                        }
                    }
            );
            if (!batch.isEmpty()) {
                int count =
                        batchWriter.write(
                                batch
                        );
                success.addAndGet(count);
                batch.clear();
            }
            taskService.finish(
                    taskId,
                    processed.get(),
                    success.get(),
                    failed.get()
            );
            publisher.complete(taskId);
        } catch (Exception e) {
            taskService.fail(
                    taskId,
                    e.getMessage()
            );
            publisher.error(
                    taskId,
                    e.getMessage()
            );
        } finally {
            try {
                Files.deleteIfExists(file);
            } catch (IOException ignored) {
            }
        }
    }
}

这样无论 Excel 是 1 万行还是 30 万行,业务层真正长期持有的 List 大概都是 500 条,而不是整个 Excel。

数据库这里我也没有再写:

for (Product product : products) {
    repository.save(product);
}

直接使用 JDBC batch。

@Repository
public class ProductBatchWriter {

    private final JdbcTemplate jdbcTemplate;
    public ProductBatchWriter(
            JdbcTemplate jdbcTemplate) {
        this.jdbcTemplate =
                jdbcTemplate;
    }
    public int write(
            List products) {
        int[][] result =
                jdbcTemplate.batchUpdate(
                        ”””
                        INSERT INTO product(
                            sku,
                            name,
                            price,
                            stock
                        )
                        VALUES (?, ?, ?, ?)
                        ”””,
                        products,
                        products.size(),
                        (ps, product) -> {

                            ps.setString(
                                    1,
                                    product.sku()
                            );

                            ps.setString(
                                    2,
                                    product.name()
                            );

                            ps.setBigDecimal(
                                    3,
                                    product.price()
                            );

                            ps.setInt(
                                    4,
                                    product.stock()
                            );
                        }
                );
        return Arrays.stream(result)
                .mapToInt(
                        batch ->
                                (int) Arrays
                                        .stream(batch)
                                        .filter(
                                                value ->
                                                        value >= 0
                                        )
                                        .count()
                )
                .sum();
    }
}

200、500、1000 都可以测。

真正不能做的是:

解析一行
INSERT 一次

10 万条就是十万次数据库交互。

这时候 Excel 解析可能根本不是瓶颈。

数据库往返才是。

另外一个以前一直让我觉得很别扭的问题,是错误数据。

假设 10 万行里面:

99,872 条正确
128 条错误

以前很多导入代码的处理方式是:

第 38127 行价格为空

throw Exception

整份 Excel 导入失败

其实没必要。

现在我把“行级错误”和“任务级错误”分开了。

比如:

SKU 为空
价格格式错误
库存小于 0
手机号格式不正确

这些属于行级问题。

这一行失败:

记录错误
继续下一行

只有:

Excel 文件损坏
数据库不可用
磁盘读取失败
表头完全不匹配

这种问题,我才让整个任务进入 FAILED。

错误数据单独写表:

CREATE TABLE import_error (
    id BIGINT PRIMARY KEY AUTO_INCREMENT,
    task_id VARCHAR(64) NOT NULL,
    row_no BIGINT NOT NULL,
    raw_data TEXT,
    error_message VARCHAR(1000),
    created_at DATETIME NOT NULL,
    INDEX idx_task_id(task_id)
);

最终任务结果可能是:

任务状态:PARTIAL_SUCCESS

总处理:100000
成功:99872
失败:128

后台再给一个“下载失败数据”,运营改这 128 条就行。

不用重新处理十万条。

做到这里其实已经能用了。

但我后来又加了一点东西,用户体验好了很多:

SSE 实时进度。

以前上传以后只能:

正在导入,请稍候……

到底还有一分钟还是十分钟,谁也不知道。

这种场景我没有上 WebSocket。

因为浏览器只需要:

服务器 → 浏览器

单向推送。

直接用 SSE 就够了。

Spring MVC 自带 SseEmitter,很适合这种“后台任务不断把进度推给前端”的场景。

我写了一个非常简单的 Publisher:

@Component
public class ImportProgressPublisher {

    private final Map<
            String,
            List
            > emitters =
            new ConcurrentHashMap<>();
    public SseEmitter subscribe(
            String taskId) {
        SseEmitter emitter =
                new SseEmitter(
                        30L * 60 * 1000
                );
        emitters.computeIfAbsent(
                taskId,
                key ->
                        new CopyOnWriteArrayList<>()
        ).add(emitter);
        emitter.onCompletion(
                () -> remove(
                        taskId,
                        emitter
                )
        );
        emitter.onTimeout(
                () -> remove(
                        taskId,
                        emitter
                )
        );
        return emitter;
    }
    public void publish(
            String taskId,
            long processed,
            long success,
            long failed) {
        ImportProgress progress =
                new ImportProgress(
                        processed,
                        success,
                        failed
                );
        List list =
                emitters.getOrDefault(
                        taskId,
                        List.of()
                );
        for (SseEmitter emitter : list) {
            try {
                emitter.send(
                        SseEmitter
                                .event()
                                .name(”progress”)
                                .data(progress)
                );
            } catch (IOException e) {
                remove(
                        taskId,
                        emitter
                );
            }
        }
    }
    public void complete(String taskId) {
        List list =
                emitters.remove(taskId);
        if (list == null) {
            return;
        }
        for (SseEmitter emitter : list) {
            try {
                emitter.send(
                        SseEmitter
                                .event()
                                .name(”complete”)
                                .data(”SUCCESS”)
                );
                emitter.complete();
            } catch (IOException ignored) {
            }
        }
    }
    public void error(
            String taskId,
            String message) {
        List list =
                emitters.remove(taskId);
        if (list == null) {
            return;
        }
        for (SseEmitter emitter : list) {
            try {
                emitter.send(
                        SseEmitter
                                .event()
                                .name(”error”)
                                .data(message)
                );
                emitter.complete();
            } catch (IOException ignored) {
            }
        }
    }
    private void remove(
            String taskId,
            SseEmitter emitter) {
        List list =
                emitters.get(taskId);
        if (list != null) {
            list.remove(emitter);
        }
    }
}

Controller:

@GetMapping(
        value = ”/api/imports/{taskId}/events”,
        produces =
                MediaType.TEXT_EVENT_STREAM_VALUE
)
public SseEmitter events(
        @PathVariable String taskId) {

    return publisher.subscribe(
            taskId
    );
}

浏览器:

const source =
    new EventSource(
        `/api/imports/${taskId}/events`
    );

source.addEventListener(
    ”progress”,
    event => {

        const data =
            JSON.parse(event.data);
        console.log(
            ”已处理:”,
            data.processed
        );

        console.log(
            ”成功:”,
            data.success
        );

        console.log(
            ”失败:”,
            data.failed
        );
    }
);

source.addEventListener(
    ”complete”,
    event => {
        console.log(”导入完成”);
        source.close();
    }
);

source.addEventListener(
    ”error”,
    event => {
        console.log(
            ”导入失败”,
            event.data
        );
        source.close();
    }
);

于是页面可以直接显示:

正在导入……

已处理:63,500
成功:63,421
失败:79

这比一个一直转圈的 Loading 好太多了。

而且 SSE 对这种场景特别合适。

因为我们根本不需要浏览器不停给服务器发送消息。

整个数据流就是:

后台任务
   ↓
Spring Boot
   ↓
SSE
   ↓
浏览器

单向推送就够了。

不过上面这个:

ConcurrentHashMap>

只适合单实例或者简单项目。

如果生产环境是:

Spring Boot A
Spring Boot B
Spring Boot C

上传请求到了 A。

SSE 连接却到了 B。

B 内存里根本没有 A 那个进度。

这时候我不会再硬改这个 Map。

真正的任务状态本来就在数据库:

processed_rows
success_rows
failed_rows
status

浏览器刷新页面以后先:

GET /api/imports/{taskId}

查询持久化状态。

实时通知如果需要跨节点,再把 progress event 放到 Redis Pub/Sub、MQ 或其他共享通道。

这样数据库负责最终可信状态,SSE 只负责让页面看起来实时。

这两个职责不要反过来。

做完以后,我还顺手解决了另外一个以前很烦的问题。

用户刷新页面。

以前刷新以后:

导入进度没了。

现在浏览器手里有 taskId。

重新请求:

GET /api/imports/{taskId}

就能拿到:

{
  ”taskId”: ”7c829...”,
  ”status”: ”RUNNING”,
  ”processedRows”: 63500,
  ”successRows”: 63421,
  ”failedRows”: 79
}

然后重新订阅 SSE。

所以现在页面刷新、关掉甚至换一台电脑,都不会影响真正的导入任务。

还有一点需要特别注意。

SSE 本身并不是任务状态存储。

如果你只把:

63500
63421
79

放在 JVM 内存里。

服务一重启,全部没了。

所以真正的状态还是要定期写数据库。

比如每处理 500 行更新一次:

UPDATE import_task
SET processed_rows = ?,
    success_rows = ?,
    failed_rows = ?
WHERE id = ?;

然后 SSE 推送同样的数据。

这意味着哪怕:

浏览器断线
Spring Boot 重启
SSE 连接失效

用户重新打开页面以后,依然可以从数据库恢复到最近一次持久化的进度。

SSE 是实时体验。

数据库才是任务事实。

回过头来看,这次改造最重要的其实不是 POI 比以前省了多少内存,也不是 500 条 batch 到底比 1000 条快多少,甚至也不只是加了 SSE。

而是我不再把:

一个需要几分钟完成的后台任务

伪装成:

一个普通 HTTP 请求。

这是很多后台系统里特别容易出现的问题。

Excel 导入。

批量生成报表。

批量发送消息。

几十万条数据同步。

批量图片处理。

AI 批量生成。

只要任务可能跑:

几十秒
几分钟
甚至十几分钟

我现在第一反应已经不是:

把 nginx timeout 调成 600 秒。

而是:

是不是应该给它一个 taskId?

HTTP 请求负责提交任务。

后台线程负责执行任务。

数据库负责保存状态。

SSE 负责实时告诉用户现在跑到哪了。

失败数据单独记录:

能成功的继续成功。

而不是第 38,127 行出了一点问题,就让前面三万多条全部陪着重来。

我现在再看最开始那段代码:

@PostMapping(”/import”)
public void importExcel(
        MultipartFile file) {

    List rows = read(file);
    repository.saveAll(rows);
}

它并不是错。

如果永远只有:

500 条

我甚至还是会这么写。

真正的问题是,我们经常拿一个原本为几百条数据设计的接口,慢慢承接:

5000
20000
50000
100000

最后出了问题以后,再不停加:

-Xmx4g
proxy_read_timeout 600
spring.mvc.async.request-timeout

但这些参数只是让问题晚一点出现。

如果一件事情本质上已经是:

后台批处理任务

最有效的优化,通常不是把 HTTP 超时时间从 30 秒改成 10 分钟。

而是别再让这个 HTTP 请求等它十分钟。

让请求快速返回一个:

taskId

让任务在后台跑。

再让 SSE 把:

已处理多少
成功多少
失败多少
是否完成

实时推给页面。

这样才是真正适合生产环境的大批量 Excel 导入。




上一篇:Xcode 27.1 beta 发布:iPhone Duo 折叠屏开发适配提前看
下一篇:数据中心半导体2031年1.53万亿:存储上位,电力成瓶颈
您需要登录后才可以回帖 登录 | 立即注册

手机版|小黑屋|网站地图|云栈社区 ( 苏ICP备2022046150号-2 )

GMT+8, 2026-9-21 02:54 , Processed in 1.133603 second(s), 41 queries , Gzip On.

Powered by Discuz! X3.5

© 2025-2026 云栈社区.

快速回复 返回顶部 返回列表