公司动态
Spring Boot 3 结合 Redis Lua 实现分布式批量数据导入内存配额控制
Spring Boot 3 结合 Redis Lua 实现分布式批量数据导入内存配额控制问题背景在企业级业务中用户、订单、商品等数据的批量导入是常见需求通常基于 EasyExcel、OpenCSV 等工具实现 Excel/CSV 文件的解析。当导入文件行数超过 10 万、或分布式集群多实例同时导入时很容易出现两类内存问题 1. 单实例导入时未控制批次大小一次性加载大量数据到内存触发 JVM OOM 2. 分布式场景下多个实例同时导入总内存占用超过集群阈值挤占其他业务的可用内存导致服务雪崩。现有方案通常存在明显短板单机限流无法控制集群总内存Redis 普通命令做配额控制时存在并发超发问题导入进度统计需要加锁或多次交互性能低下。本文将结合 Spring Boot 3 的业务编排能力和 Redis Lua 脚本的原子执行特性实现一套分布式批量导入的内存配额控制方案。方案设计技术分工Spring Boot 3作为核心业务框架负责文件解析、批次调度、业务逻辑处理、接口暴露不直接参与内存配额的计算和状态存储仅通过 RedisTemplate 调用 Lua 脚本完成原子操作。Redis Lua 脚本负责两类核心原子操作一是分布式内存配额令牌桶的扣减、补充、退还保证多实例同时导入时总内存配额不超限二是导入进度的原子更新避免并发更新导致的数据不一致。整体流程用户上传 Excel/CSV 文件后端生成唯一任务 ID初始化内存配额如总配额 1024MB每秒补充 100MB和进度统计结构EasyExcel 逐行读取数据每积累 1000 行1 批调用 Lua 脚本扣减对应内存令牌每行预估 1KB1 批扣减 1MB令牌扣减成功则执行业务逻辑校验、入库处理完成后调用 Lua 脚本原子更新进度令牌不足时等待 1 秒后重试直到令牌补充足够处理异常时退还已扣减的令牌避免配额无效占用前端通过任务 ID 轮询导入进度。关键原理1. 内存配额令牌桶采用令牌桶算法控制内存配额相比固定窗口算法更适合导入场景的内存波动特性令牌桶允许一定程度的突发内存占用比如单批处理需要较高内存同时长期平均内存占用不超过设定的总配额。令牌桶状态存储在 Redis Hash 结构中包含当前令牌数、上次补充时间、桶容量、 replenish 速率四个字段所有状态更新由 Lua 脚本原子执行避免多实例并发修改导致的状态不一致。2. Lua 脚本的原子性保证Redis 采用单线程模型处理命令Lua 脚本在执行过程中会独占 Redis 线程不会被其他命令插队因此多个实例同时调用配额扣减脚本时不会出现超发问题保证集群总内存配额永远不会超过设定的阈值。同时进度更新脚本利用 Redis 的HINCRBY原子命令无需加锁即可实现多实例同时更新同一个任务的进度。3. Spring Boot 3 的集成优化Spring Boot 3 的 Spring Data Redis 对 Lua 脚本的执行做了优化支持预加载脚本获取 SHA1 值后续调用时通过EVALSHA执行减少网络传输开销提升高并发下的调用性能。同时 Spring Boot 3 的异步任务处理能力避免导入任务阻塞主线程提升接口响应速度。完整示例环境要求JDK 17Spring Boot 3.2.xEasyExcel 3.3.xRedis 6.x1. 依赖引入dependencies dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependency dependency groupIdcom.alibaba/groupId artifactIdeasyexcel/artifactId version3.3.2/version /dependency dependency groupIdorg.projectlombok/groupId artifactIdlombok/artifactId optionaltrue/optional /dependency /dependencies2. Lua 脚本定义与预加载Component public class LuaScriptHolder { // 内存配额扣减Lua脚本 public static final String MEMORY_QUOTA_SCRIPT local quota_key KEYS[1] local need_token tonumber(ARGV[1]) local now redis.call(TIME)[1] local current_token tonumber(redis.call(HGET, quota_key, current_token) or 0) local last_replenish tonumber(redis.call(HGET, quota_key, last_replenish_time) or now) local capacity tonumber(redis.call(HGET, quota_key, capacity) or 1024) local rate tonumber(redis.call(HGET, quota_key, replenish_rate) or 100) local past_time now - last_replenish if past_time 0 then local replenish past_time * rate current_token math.min(current_token replenish, capacity) end if current_token need_token then current_token current_token - need_token redis.call(HSET, quota_key, current_token, current_token, last_replenish_time, now) return 1 else redis.call(HSET, quota_key, last_replenish_time, now) return 0 end; // 导入进度更新Lua脚本 public static final String PROGRESS_UPDATE_SCRIPT local progress_key KEYS[1] local batch_rows tonumber(ARGV[1]) local success_rows tonumber(ARGV[2]) local fail_rows tonumber(ARGV[3]) redis.call(HINCRBY, progress_key, total_rows, batch_rows) redis.call(HINCRBY, progress_key, success_rows, success_rows) redis.call(HINCRBY, progress_key, fail_rows, fail_rows) redis.call(HSET, progress_key, update_time, redis.call(TIME)[1]) return redis.call(HGETALL, progress_key); Autowired private RedisTemplateString, Object redisTemplate; private String memoryQuotaScriptSha; private String progressUpdateScriptSha; PostConstruct public void init() { // 预加载脚本到Redis获取SHA1值后续调用性能更高 memoryQuotaScriptSha redisTemplate.opsForValue().getScriptingCommands().scriptLoad(MEMORY_QUOTA_SCRIPT); progressUpdateScriptSha redisTemplate.opsForValue().getScriptingCommands().scriptLoad(PROGRESS_UPDATE_SCRIPT); } public String getMemoryQuotaScriptSha() { return memoryQuotaScriptSha; } public String getProgressUpdateScriptSha() { return progressUpdateScriptSha; } }3. EasyExcel 监听器实现Component Slf4j public class ImportDataListener extends AnalysisEventListenerMapInteger, String { private static final int BATCH_SIZE 1000; // 每批处理行数 private static final int PER_ROW_MEMORY_KB 1; // 每行预估内存KB根据业务调整 private ListMapInteger, String batchData new ArrayList(BATCH_SIZE); private String taskId; private final RedisTemplateString, Object redisTemplate; private final LuaScriptHolder luaScriptHolder; public ImportDataListener(RedisTemplateString, Object redisTemplate, LuaScriptHolder luaScriptHolder) { this.redisTemplate redisTemplate; this.luaScriptHolder luaScriptHolder; } public void setTaskId(String taskId) { this.taskId taskId; } Override public void invoke(MapInteger, String data, AnalysisContext context) { batchData.add(data); if (batchData.size() BATCH_SIZE) { processBatch(); } } private void processBatch() { if (batchData.isEmpty()) return; // 计算本批需要扣减的令牌MB int needToken (BATCH_SIZE * PER_ROW_MEMORY_KB) / 1024; // 调用Lua脚本扣减配额 Long result redisTemplate.execute( new DefaultRedisScript() {{ setScriptText(luaScriptHolder.getMemoryQuotaScriptSha()); setResultType(Long.class); }}, Collections.singletonList(import:memory:quota: taskId), String.valueOf(needToken) ); if (result ! null result 1) { int success 0, fail 0; try { for (MapInteger, String row : batchData) { if (validateRow(row) saveRow(row)) { success; } else { fail; } } } catch (Exception e) { log.error(导入批次处理失败任务ID{}, taskId, e); fail batchData.size(); refundToken(needToken); // 处理失败退还令牌 } finally { // 原子更新进度 redisTemplate.execute( new DefaultRedisScript() {{ setScriptText(luaScriptHolder.getProgressUpdateScriptSha()); setResultType(List.class); }}, Collections.singletonList(import:progress: taskId), String.valueOf(batchData.size()), String.valueOf(success), String.valueOf(fail) ); batchData.clear(); } } else { // 令牌不足等待后重试 log.warn(任务{}内存配额不足等待补充, taskId); try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } processBatch(); } } // 退还令牌 private void refundToken(int token) { redisTemplate.opsForHash().increment(import:memory:quota: taskId, current_token, token); } private boolean validateRow(MapInteger, String row) { // 业务校验逻辑 return true; } private boolean saveRow(MapInteger, String row) { // 业务入库逻辑 return true; } Override public void doAfterAllAnalysed(AnalysisContext context) { processBatch(); redisTemplate.opsForHash().put(import:progress: taskId, status, completed); } Override public void onException(Exception exception, AnalysisContext context) { log.error(Excel解析异常任务ID{}, taskId, exception); if (!batchData.isEmpty()) { int needToken (batchData.size() * PER_ROW_MEMORY_KB) / 1024; refundToken(needToken); batchData.clear(); } } }4. 导入控制器与异步工具RestController RequestMapping(/api/import) public class ImportController { private final ImportDataListener importDataListener; private final RedisTemplateString, Object redisTemplate; public ImportController(ImportDataListener importDataListener, RedisTemplateString, Object redisTemplate) { this.importDataListener importDataListener; this.redisTemplate redisTemplate; } PostMapping(/excel) public MapString, String importExcel(RequestParam(file) MultipartFile file) throws IOException { String taskId UUID.randomUUID().toString(); // 初始化配额总容量1024MB每秒补充100MB initQuota(taskId, 1024, 100); // 初始化进度 initProgress(taskId); importDataListener.setTaskId(taskId); // 异步执行导入 AsyncManager.execute(() - { try { EasyExcel.read(file.getInputStream(), importDataListener).sheet().doRead(); } catch (IOException e) { log.error(导入失败任务ID{}, taskId, e); redisTemplate.opsForHash().put(import:progress: taskId, status, failed); } }); return Map.of(taskId, taskId); } private void initQuota(String taskId, int capacityMB, int ratePerSec) { String key import:memory:quota: taskId; redisTemplate.opsForHash().putAll(key, Map.of( current_token, String.valueOf(capacityMB), last_replenish_time, String.valueOf(System.currentTimeMillis() / 1000), capacity, String.valueOf(capacityMB), replenish_rate, String.valueOf(ratePerSec) )); } private void initProgress(String taskId) { String key import:progress: taskId; redisTemplate.opsForHash().putAll(key, Map.of( total_rows, 0, success_rows, 0, fail_rows, 0, update_time, String.valueOf(System.currentTimeMillis() / 1000), status, running )); } GetMapping(/progress/{taskId}) public MapString, Object getProgress(PathVariable String taskId) { return redisTemplate.opsForHash().entries(import:progress: taskId); } }Component public class AsyncManager { private static final ExecutorService EXECUTOR new ThreadPoolExecutor( 5, 10, 60L, TimeUnit.SECONDS, new LinkedBlockingQueue(100), new ThreadFactory() { private final AtomicInteger count new AtomicInteger(1); Override public Thread newThread(Runnable r) { return new Thread(r, import-executor- count.getAndIncrement()); } }, new ThreadPoolExecutor.CallerRunsPolicy() ); public static void execute(Runnable task) { EXECUTOR.execute(task); } }常见问题1. Lua 脚本执行超时Redis 单线程处理命令若 Lua 脚本包含复杂循环或阻塞操作会阻塞后续所有命令。本方案的两个脚本均为简单的哈希操作执行时间在 1ms 以内不会出现超时问题。若业务需要更复杂的逻辑建议拆分为多个短脚本避免单脚本执行时间过长。2. 配额控制精度不足本方案采用预估每行内存的方式扣减令牌若业务数据字段差异大、或存在复杂对象处理预估值和实际内存占用可能存在偏差。可通过压测调整PER_ROW_MEMORY_KB参数或根据实际业务逻辑动态计算每批的内存占用提升控制精度。3. 令牌退还遗漏若处理批次时发生异常未退还令牌会导致配额被无效占用后续导入任务无法获取足够令牌。本方案在catch块、finally块和onException方法中均加入了退还逻辑保证异常场景下令牌正确释放。4. 分布式时间不同步若使用 Java 系统时间计算令牌补充量Java 服务器与 Redis 服务器时间不同步会导致配额计算错误。本方案 Lua 脚本中使用 Redis 原生的TIME命令获取服务器时间避免时间偏差问题。适用边界与关键取舍适用场景分布式部署的批量导入场景单文件行数超过 10 万需要控制集群总内存占用需要对导入进度做统一统计、支持多实例同时导入的场景。不适用场景单机部署、导入文件小于 1 万行的轻量场景无需引入 Redis使用本地令牌桶即可单批处理内存占用超过总配额的场景如单批需要处理 500MB 数据总配额仅 512MB令牌桶控制意义有限需调整批次大小或采用其他方案。关键取舍令牌桶 vs 固定窗口令牌桶允许突发内存占用适合导入场景的内存波动特性但控制精度略低于固定窗口若业务对内存控制要求极高、不允许任何突发可替换为固定窗口算法。Lua 脚本 vs 普通 Redis 命令Lua 脚本保证原子性避免并发超发但脚本长度不能过长否则阻塞 Redis若并发量极低也可用WATCHMULTI事务实现但性能低于 Lua 脚本。预估内存 vs 实际内存统计预估内存实现简单、性能高但精度有限若需要精确控制可在批次处理完成后通过Runtime.getRuntime().totalMemory() - Runtime.freeMemory()计算实际内存占用再调整令牌但会增加性能开销。容易踩坑的细节Lua 脚本预加载不要每次调用都传递完整脚本内容预加载脚本获取 SHA1 值后通过EVALSHA调用可减少 50% 以上的网络传输开销高并发下性能提升明显。批次大小与配额粒度匹配批次大小不要超过配额粒度的 10 倍比如配额粒度是 1MB批次大小不要超过 10MB避免单批内存占用超过总配额导致导入失败。Redis 持久化配置若 Redis 宕机导致令牌桶状态丢失配额控制会失效建议开启 AOF 持久化并设置合理的刷盘策略或定期将配额状态落盘到数据库宕机后恢复。总结本方案通过 Spring Boot 3 实现业务逻辑的编排利用 Redis Lua 脚本的原子执行特性解决了分布式批量导入场景下的内存配额控制和进度统计问题既保证了多实例导入时集群总内存不超限又避免了并发场景下的数据不一致问题可支持单文件百万行级别的批量导入同时将集群内存占用控制在预设阈值内。