问题描述:

项目上需要导入十万级别的excel数据,刚开始写没考虑到数据量会这么大,导致导入数据的时候应用程序容易报内存溢出


原因分析

  1. 导入需要集合来存储 Excel 行列的值,当这个批次被导入完成后,集合没进行清空,一直占用着内存
  2. 数据插入到数据的时候使用的是MyBatis 的批量插入,MyBatis 在处理结果集的时候是一条一条数据进行循环遍历处理的,效率明显不高,MyBatis 底层是通过 JDBC 对数据库进行操作的,所以我们可以直接使用JDBC进行批量操作

解决方案:

针对上面分析的两点原因,作出以下改造

EasyExcel 监听器

package com.southsmart.sgeocserver.analyze.pipe.listener;

import com.alibaba.excel.annotation.ExcelProperty;
import com.alibaba.excel.context.AnalysisContext;
import com.alibaba.excel.event.AnalysisEventListener;
import com.alibaba.excel.exception.ExcelAnalysisException;
import lombok.SneakyThrows;
import org.springframework.util.StringUtils;

import java.lang.reflect.Field;
import java.util.*;

/**
 * The type Point read listener.
 *
 * @version 1.0
 */
public class PointReadListener extends AnalysisEventListener<ExcelPointVO> {

    private final Class<ExcelPointVO> tClass;

    protected final IPointService pointService;

    public PointReadListener(Class<ExcelPointVO> tClass, IPointService pointService) {
        this.tClass = tClass;
        this.pointService= pointService;
    }
    
    /**
     * The List.
     */
    List<ExcelPointVO> list = new ArrayList<>();

    @Override
    public void invokeHeadMap(Map<Integer, String> headMap, AnalysisContext context) {
        super.invokeHeadMap(headMap, context);
    }

    /**
     * All listeners receive this method when any one Listener does an error report. If an exception is thrown here, the
     * entire read will terminate.
     *
     * @param exception
     * @param context
     */
    @Override
    public void onException(Exception exception, AnalysisContext context) throws Exception {
        list.clear();
        super.onException(exception, context);
    }

    /**
     * When analysis one row trigger invoke function.
     *
     * @param data    one row value. Is is same as {@link AnalysisContext#readRowHolder()}
     * @param context
     */
    @SneakyThrows
    @Override
    public void invoke(ExcelPointVO data, AnalysisContext context) {
        Class<? extends ExcelPointVO> clazz = data.getClass();
        Field[] declaredFields = clazz.getDeclaredFields();
        for (Field declaredField : declaredFields) {
            declaredField.setAccessible(true);
            Class<?> type = declaredField.getType();
            if (type.isAssignableFrom(String.class) && !StringUtils.hasText((String) declaredField.get(data))) {
                declaredField.set(data, "");
            }
        }
        list.add(data);
        if (list.size() >= Constants.GENERAL_ONCE_SAVE_TO_DB_ROWS_JDBC) {
            // 存入数据库:数据小于 2000 条使用批量插入即可
            saveData();
            // 清理集合便于GC回收
            list.clear();
        }
    }

    /**
     * 保存数据到 DB
     */
    private void saveData() {
        if (list.size() > 0) {
            pointService.batchSaveImportPointData(list);
            list.clear();
        }
    }

    /**
     * if have something to do after all analysis
     *
     * @param context c
     */
    @Override
    public void doAfterAllAnalysed(AnalysisContext context) {
        saveData();
        list.clear();
    }
}

JDBC批量插入实现类

    public Map<String, Object> batchSaveImportPointData(List<ExcelPointVO> dataList) {
        // 插入数据库
        Map<String, List<ExcelPointVO>> listMap = dataList.stream().collect(Collectors.groupingBy(ExcelPointVO::getTableType));
        Map<String, Object> result = new HashMap<>();
        // 结果集中数据为0时,结束方法.进行下一次调用
        if (dataList.size() == 0) {
            return result ;
        }
        for (Map.Entry<String, List<ExcelPointVO>> entry : listMap.entrySet()) {
            // JDBC分批插入+事务操作完成对2000数据的插入
            Connection conn = null;
            PreparedStatement ps1 = null;
            try {
                long startTime = System.currentTimeMillis();
                // 获取表名称
                String name = "你的表名";
                // 插入数据
                StringBuilder stringBuilder = new StringBuilder();
                //处理之后的业务sql
                String sql = dealSql(entry.getValue());
                stringBuilder.append("INSERT INTO ").append(name).append("(\"id\", \"modelid\", ...) VALUES ")
                        .append(" ").append(sql).append(" ON conflict(id) DO NOTHING ");
                log.info("{} 条,开始导入到数据库时间:{}", entry.getValue().size(), startTime + "ms");
                conn = jdbcDruidUtils.getConnection();
                // 控制事务:默认不提交
                conn.setAutoCommit(false);
                ps1 = conn.prepareStatement(stringBuilder.toString());
                // 执行批处理
                ps1.execute();
                // 手动提交事务
                conn.commit();
                long endTime = System.currentTimeMillis();
                log.info("{} 条,结束导入到数据库时间:{}", entry.getValue().size(), endTime + "ms");
                log.info("{} 条,导入用时:{}", entry.getValue().size(), (endTime - startTime) + "ms");
                result.put("success", "1");
            } catch (Exception e) {
                result.put("exception", "0");
                e.printStackTrace();
            } finally {
                // 关连接
                JDBCDruidUtils.close(ps1);
                JDBCDruidUtils.close(conn);
            }
        }
        return result;
    }

JDBCDruidUtils工具类

package com.southsmart.sgeocserver.analyze.util;

import org.springframework.stereotype.Component;

import javax.annotation.Resource;
import javax.sql.DataSource;
import java.sql.Connection;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Statement;

@Component
public class JDBCDruidUtils {
    /**
     * 创建连接池对象
     */
    @Resource
    private DataSource dataSource;

    /**
     * 从连接池里获取一个连接对象
     *
     * @return Connection接口
     * @throws SQLException 以后使用的时候需要知道这个方法是否出现异常,用于警告,所以抛出
     */
    public Connection getConnection() throws SQLException {
        return dataSource.getConnection();
    }

    /**
     * 释放(归还)资源
     *
     * @param stmt Statement 执行SQL的对象
     * @param conn Connection 数据库连接对象
     */
    public static void close(Statement stmt, Connection conn) {
        //调用三个参数的close,第一个参数传null
        close(null, stmt, conn);
    }

    /**
     * 释放(归还)资源
     *
     * @param stmt Statement 执行SQL的对象
     */
    public static void close(Statement stmt) {
        //调用三个参数的close,第一个参数传null
        close(null, stmt, null);
    }

    /**
     * 释放(归还)资源
     *
     * @param conn Connection 数据库连接对象
     */
    public static void close(Connection conn) {
        //调用三个参数的close,第一个参数传null
        close(null, null, conn);
    }

    /**
     * 释放(归还)资源
     *
     * @param rs   ResultSet 结果集
     * @param stmt Statement 执行SQL的对象
     * @param conn Connection 数据库连接对象
     */
    public static void close(ResultSet rs, Statement stmt, Connection conn) {
        if (rs != null) {
            try {
                rs.close();
            } catch (SQLException e) {
                e.printStackTrace();
            }
        }

        if (stmt != null) {
            try {
                stmt.close();
            } catch (SQLException e) {
                e.printStackTrace();
            }
        }

        if (conn != null) {
            try {
                conn.close();
            } catch (SQLException e) {
                e.printStackTrace();
            }
        }
    }
}

看了其他人的优化方案,后面可以考虑再加上多线程处理进行改造

参考链接

  1. https://blog.csdn.net/weixin_36754290/article/details/123715039
  2. https://blog.csdn.net/vnjohn/article/details/125980327
Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐