尧图精选

超长视频上传与动态定价:高并发场景下的Java工程实践

🕒 发布时间:2026/9/2 21:55:23 📁 来源:尧图网络
最近在做一个视频平台的需求时被一个业务提法折腾了好几个晚上“超长视频来袭午高峰两小时恶劣天气单价这么低”乍看像一句运营吐槽实际上拆开全是技术问题超长视频怎么稳定上传午高峰两小时流量怎么扛恶劣天气下动态价格怎么算才合理这篇文章从这几个问题出发整理了一套从分片上传、异步转码、限流熔断到动态定价规则的完整落地思路包含可复制的 Java 代码和实践建议。1. 业务场景与技术痛点拆解1.1 超长视频带来的三个直接挑战先理解“超长视频”在技术上的含义。以前平台常见视频是几分钟到十几分钟的短视频单文件几十 MB上传失败重传成本不高。但长视频、课程回放、监控录像、直播录制单文件动辄几 GB甚至几十 GB直接上传会暴露出一连串问题网络波动导致上传中断用户只能重新开始体验极差。一次性上传大文件服务端接收耗时久容易触发网关超时。视频转码耗时长同步接口很容易把请求线程池占满。所以超长视频首先需要解决的不是播放而是“怎么把一个大文件稳定收上来”。分片上传和断点续传是这个场景的基础能力。1.2 午高峰两小时意味着什么视频平台的流量并不均匀午休时间通常 12:00-14:00是上传和播放的双高峰。这个时段的特点是上传并发量突增大量用户同时提交大视频。播放请求也随之上涨转码、分发、带宽成本同步升高。如果系统按峰值时刻的 2 倍设计资源浪费严重如果不预留高峰期又会抖动。午高峰给技术团队的启示是系统必须具备削峰填谷能力不能指望靠加机器硬扛。更好的方案是异步化、队列缓冲、限流降级、弹性扩容组合使用。1.3 “恶劣天气单价这么低”到底在说什么这句话看起来像运营对定价策略的抱怨实际上是一个动态定价问题。很多服务型平台会根据天气调整价格比如雨天网约车溢价、恶劣天气外卖配送费上涨。但这里的场景是“恶劣天气单价低”也就是业务侧觉得价格策略不够合理。从技术角度看动态定价需要一个规则引擎获取天气数据、匹配计费规则、计算价格、灰度发布、监控效果。单价低不低不是拍脑袋决定的而是由供需关系、天气影响系数、历史成交数据共同计算出来的。本文不会讨论具体的商业策略而是从工程实现角度给出“天气数据接入 动态定价规则引擎”的最小实现方便你在业务侧做策略调整时快速落地。2. 环境准备与版本说明为了保证示例可运行本文采用如下环境。版本需要根据你的项目实际情况调整本文重点演示配置思路。组件版本/说明操作系统Windows / macOS / Linux 均可JDK8 及以上建议 17Spring Boot2.7.x 或 3.xRedis5.0 及以上用于分片状态记录与限流RabbitMQ3.x用于异步任务解耦MySQL5.7 及以上存储视频信息和价格规则FFmpeg4.x用于视频转码MinIO用于本地模拟对象存储生产可替换为 OSS/COS如果你本机没有 RabbitMQ 和 MinIO也可以用 Docker 快速启动docker run -d --name redis -p 6379:6379 redis:7 docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:3-management docker run -d --name minio -p 9000:9000 -p 9001:9001 minio/minio server /data --console-address :9001建议新建一个统一的 Maven 工程示例项目名称为video-upload-demo。3. 核心原理拆解3.1 分片上传与断点续传原理分片上传的核心思想很简单把一个大文件切分成多个小分片逐个上传服务端记录每个分片的上传状态全部上传完成后由服务端触发合并。分片上传解决了以下问题网络中断后只需要重传未完成的分片而不是整个文件。单个请求体积可控网关超时概率大幅降低。支持并发上传多个分片加快整体上传速度。断点续传的关键在于“分片状态记录”。常见做法是以前端生成的文件唯一标识如 MD5作为 key在 Redis 中记录已完成的分片编号。前端每次上传前先查询已上传分片只上传缺失部分。3.2 异步转码与消息队列解耦视频上传完成后通常需要转码成多种清晰度格式例如 720P、1080P。转码是 CPU 密集型任务耗时可能超过上传本身。如果使用同步接口用户请求会长时间阻塞非常不友好。正确做法是上传合并完成后发送一条消息到 RabbitMQ由独立消费者执行转码任务。这样上传接口可以快速响应转码过程对用户透明。用户通过轮询或回调获取转码进度。3.3 高峰流量下的限流与熔断午高峰两小时流量是平时的数倍。限流不是限制用户而是保护下游核心服务不被瞬时流量打垮。常用的限流算法有固定窗口、滑动窗口、令牌桶、漏桶。在 Spring Boot 项目中可以使用 Resilience4j 或 Sentinel。Sentinel 的流量控制能力更全面还支持熔断降级。个人推荐在视频上传场景中使用 Sentinel因为上传类接口的流量往往非常集中。3.4 动态定价规则引擎动态定价通常是这样的流程收集基础价格、查询天气和供需系数、按规则计算最终价格。规则需要可配置不能在代码里写死。一个轻量级做法是把价格规则配置在数据库表中通过配置中心或定时刷新加载规则引擎根据条件依次匹配。天气数据可以来自第三方 API也可以来自自建的气象数据同步服务。4. 完整实战案例下面我们实现一个简化但完整的版本。功能包括分片上传接口。分片合并与转码消息发送。午高峰限流配置。天气数据接入与动态定价规则计算。4.1 创建项目结构video-upload-demo ├── pom.xml ├── src/main/java/com/example/videodemo │ ├── VideoDemoApplication.java │ ├── controller │ │ ├── UploadController.java │ │ └── PriceController.java │ ├── service │ │ ├── UploadService.java │ │ ├── TranscodeService.java │ │ ├── WeatherService.java │ │ └── PriceRuleEngine.java │ ├── entity │ │ └── VideoFile.java │ ├── config │ │ ├── RedisConfig.java │ │ └── RabbitConfig.java │ └── common │ └── Result.java └── src/main/resources ├── application.yml └── price-rules.json4.2 pom.xml 核心依赖dependencies dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency dependency groupIdorg.mybatis.spring.boot/groupId artifactIdmybatis-spring-boot-starter/artifactId version2.3.2/version /dependency dependency groupIdcom.alibaba.cloud/groupId artifactIdspring-cloud-starter-alibaba-sentinel/artifactId version2021.0.5.0/version /dependency /dependencies如果 Spring Boot 版本较高建议使用与当前版本匹配的 cloud 版本避免依赖冲突。4.3 application.yml 配置server: port: 8080 spring: redis: host: localhost port: 6379 rabbitmq: host: localhost port: 5672 username: guest password: guest servlet: multipart: max-file-size: 100MB max-request-size: 1GB minio: endpoint: http://localhost:9000 access-key: minioadmin secret-key: minioadmin bucket: video-bucket需要注意Spring 的 multipart 配置限制了单次请求大小。分片上传中每个分片大小一般设置为 5MB 到 10MB所以 max-file-size 配置为 100MB 足够覆盖大多数情况。4.4 分片上传接口先定义统一返回类和上传实体类。// 文件路径src/main/java/com/example/videodemo/common/Result.java package com.example.videodemo.common; public class ResultT { private int code; private String message; private T data; public static T ResultT success(T data) { ResultT result new Result(); result.code 0; result.message success; result.data data; return result; } public static T ResultT error(String message) { ResultT result new Result(); result.code 500; result.message message; return result; } public int getCode() { return code; } public void setCode(int code) { this.code code; } public String getMessage() { return message; } public void setMessage(String message) { this.message message; } public T getData() { return data; } public void setData(T data) { this.data data; } }// 文件路径src/main/java/com/example/videodemo/entity/VideoFile.java package com.example.videodemo.entity; public class VideoFile { private String fileId; private String fileName; private Long totalSize; private Integer totalChunks; private String uploadStatus; private String storagePath; // getter / setter 省略使用 IDE 自动生成 }接下来是上传 Controller 和 Service。// 文件路径src/main/java/com/example/videodemo/controller/UploadController.java package com.example.videodemo.controller; import com.example.videodemo.common.Result; import com.example.videodemo.service.UploadService; import org.springframework.web.bind.annotation.*; import org.springframework.web.multipart.MultipartFile; import javax.annotation.Resource; RestController RequestMapping(/upload) public class UploadController { Resource private UploadService uploadService; PostMapping(/chunk) public ResultString uploadChunk(RequestParam(fileId) String fileId, RequestParam(chunkIndex) Integer chunkIndex, RequestParam(totalChunks) Integer totalChunks, RequestParam(file) MultipartFile file) { uploadService.uploadChunk(fileId, chunkIndex, totalChunks, file); return Result.success(分片上传成功); } PostMapping(/merge) public ResultString merge(RequestParam(fileId) String fileId, RequestParam(fileName) String fileName) { String storagePath uploadService.mergeChunks(fileId, fileName); return Result.success(storagePath); } GetMapping(/progress) public ResultInteger getProgress(RequestParam(fileId) String fileId) { return Result.success(uploadService.getUploadProgress(fileId)); } }UploadService 负责分片的读写、状态记录和合并。// 文件路径src/main/java/com/example/videodemo/service/UploadService.java package com.example.videodemo.service; import org.springframework.data.redis.core.RedisTemplate; import org.springframework.stereotype.Service; import org.springframework.web.multipart.MultipartFile; import javax.annotation.Resource; import java.io.File; import java.io.FileOutputStream; import java.io.InputStream; import java.nio.file.Files; import java.nio.file.Path; import java.nio.file.Paths; import java.util.Set; Service public class UploadService { private static final String UPLOAD_DIR /data/video-tmp/; private static final String FILE_KEY_PREFIX upload:file:; Resource private RedisTemplateString, Object redisTemplate; Resource private TranscodeService transcodeService; public void uploadChunk(String fileId, Integer chunkIndex, Integer totalChunks, MultipartFile file) { try { Path chunkDir Paths.get(UPLOAD_DIR, fileId); if (!Files.exists(chunkDir)) { Files.createDirectories(chunkDir); } Path chunkFile chunkDir.resolve(chunkIndex .part); try (InputStream in file.getInputStream(); FileOutputStream out new FileOutputStream(chunkFile.toFile())) { byte[] buffer new byte[1024 * 1024]; int len; while ((len in.read(buffer)) ! -1) { out.write(buffer, 0, len); } } redisTemplate.opsForSet().add(FILE_KEY_PREFIX fileId, chunkIndex); redisTemplate.opsForValue().set(FILE_KEY_PREFIX fileId :total, totalChunks); } catch (Exception e) { throw new RuntimeException(分片写入失败, e); } } public String mergeChunks(String fileId, String fileName) { Integer totalChunks (Integer) redisTemplate.opsForValue().get(FILE_KEY_PREFIX fileId :total); if (totalChunks null) { throw new RuntimeException(分片信息不存在); } Path chunkDir Paths.get(UPLOAD_DIR, fileId); Path targetFile Paths.get(UPLOAD_DIR, fileId _ fileName); try (FileOutputStream out new FileOutputStream(targetFile.toFile())) { for (int i 0; i totalChunks; i) { Path part chunkDir.resolve(i .part); if (!Files.exists(part)) { throw new RuntimeException(分片缺失: i); } Files.copy(part, out); } } catch (Exception e) { throw new RuntimeException(合并分片失败, e); } // 清理临时分片目录 deleteDir(chunkDir.toFile()); redisTemplate.delete(FILE_KEY_PREFIX fileId); redisTemplate.delete(FILE_KEY_PREFIX fileId :total); // 发送转码消息 transcodeService.sendTranscodeTask(targetFile.toString()); return targetFile.toString(); } public Integer getUploadProgress(String fileId) { SetObject uploaded redisTemplate.opsForSet().members(FILE_KEY_PREFIX fileId); Integer total (Integer) redisTemplate.opsForValue().get(FILE_KEY_PREFIX fileId :total); if (uploaded null || total null || total 0) { return 0; } return uploaded.size() * 100 / total; } private void deleteDir(File dir) { File[] files dir.listFiles(); if (files ! null) { for (File f : files) { f.delete(); } } dir.delete(); } }这段代码使用了 Redis 的 Set 结构记录已上传分片用 String 结构记录总分片数。查询进度时用已上传分片数除以总分片数简单可靠。4.5 异步转码消息上传合并完成后不应该让用户等待转码完成。这里通过 RabbitMQ 解耦。// 文件路径src/main/java/com/example/videodemo/config/RabbitConfig.java package com.example.videodemo.config; import org.springframework.amqp.core.Queue; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; Configuration public class RabbitConfig { public static final String TRANSCODE_QUEUE video.transcode.queue; Bean public Queue transcodeQueue() { return new Queue(TRANSCODE_QUEUE, true); } }// 文件路径src/main/java/com/example/videodemo/service/TranscodeService.java package com.example.videodemo.service; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.stereotype.Service; import javax.annotation.Resource; Service public class TranscodeService { Resource private RabbitTemplate rabbitTemplate; public void sendTranscodeTask(String storagePath) { rabbitTemplate.convertAndSend(RabbitConfig.TRANSCODE_QUEUE, storagePath); } }消费者接收消息后调用 FFmpeg 命令执行转码。因为 FFmpeg 命令是外部进程建议使用 ProcessBuilder 调用。// 文件路径src/main/java/com/example/videodemo/service/TranscodeConsumer.java package com.example.videodemo.service; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Component; import java.io.BufferedReader; import java.io.InputStreamReader; Component public class TranscodeConsumer { RabbitListener(queues RabbitConfig.TRANSCODE_QUEUE) public void handleTranscode(String storagePath) { try { ProcessBuilder pb new ProcessBuilder( ffmpeg, -i, storagePath, -b:v, 1M, -c:v, libx264, storagePath _720p.mp4 ); Process process pb.start(); int exitCode process.waitFor(); if (exitCode 0) { System.out.println(转码成功: storagePath); } else { System.err.println(转码失败: storagePath); } } catch (Exception e) { System.err.println(转码异常: e.getMessage()); } } }这里使用-b:v 1M是一个示例参数实际生产环境需要根据视频分辨率和清晰度要求做更细的编码参数配置。4.6 午高峰限流配置午高峰两小时流量大我们可以对上传接口做限流。使用 Sentinel 的注解方式非常方便。// 文件路径src/main/java/com/example/videodemo/controller/UploadController.java // 只展示新增部分 PostMapping(/chunk) SentinelResource(value uploadChunk, blockHandler blockHandlerForUpload) public ResultString uploadChunk(...) { // 原逻辑不变 } public ResultString blockHandlerForUpload(...) { return Result.error(当前上传人数太多请稍后再试); }在启动类或配置类中添加 Sentinel 规则初始化// 文件路径src/main/java/com/example/videodemo/config/SentinelConfig.java package com.example.videodemo.config; import com.alibaba.csp.sentinel.slots.block.RuleConstant; import com.alibaba.csp.sentinel.slots.block.flow.FlowRule; import com.alibaba.csp.sentinel.slots.block.flow.FlowRuleManager; import org.springframework.context.annotation.Configuration; import javax.annotation.PostConstruct; import java.util.ArrayList; import java.util.List; Configuration public class SentinelConfig { PostConstruct public void init() { ListFlowRule rules new ArrayList(); FlowRule rule new FlowRule(); rule.setResource(uploadChunk); rule.setGrade(RuleConstant.FLOW_GRADE_QPS); rule.setCount(1000); rules.add(rule); FlowRuleManager.loadRules(rules); } }限流的 QPS 阈值需要根据线上压测结果调整。午高峰时如果阈值设置过低大量用户会被限流设置过高又起不到保护作用。建议最初的阈值设置为正常峰值的 1.5 倍再逐步调优。4.7 天气数据接入与动态定价动态定价是本例的另一个重点。我们可以用两个组件实现一个负责获取天气一个负责计算价格。首先定义一个天气服务这里使用第三方公开天气 API 的通用查询模式具体接口以你选择的供应商为准。// 文件路径src/main/java/com/example/videodemo/service/WeatherService.java package com.example.videodemo.service; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import org.springframework.web.client.RestTemplate; import java.util.Map; Service public class WeatherService { Value(${weather.api-url}) private String apiUrl; Value(${weather.api-key}) private String apiKey; private final RestTemplate restTemplate new RestTemplate(); public String getWeatherByCity(String city) { String url apiUrl ?city city key apiKey; MapString, Object response restTemplate.getForObject(url, Map.class); if (response ! null response.get(results) ! null) { // 这里简化解析实际需要根据返回结构提取天气类型 return rain; } return unknown; } }然后定义一个价格规则引擎。规则使用配置中心或本地 JSON 文件管理方便运营修改。# application.yml 增加配置 weather: api-url: https://api.example.com/weather api-key: your-key规则文件// 文件路径src/main/resources/price-rules.json { basePrice: 10, rules: [ { condition: rain, supplyFactor: 0.8, demandFactor: 1.5 }, { condition: snow, supplyFactor: 0.7, demandFactor: 2.0 }, { condition: sunny, supplyFactor: 1.0, demandFactor: 1.0 } ] }价格引擎读取规则计算最终价格。// 文件路径src/main/java/com/example/videodemo/service/PriceRuleEngine.java package com.example.videodemo.service; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.beans.factory.annotation.Value; import org.springframework.core.io.Resource; import org.springframework.stereotype.Service; import javax.annotation.PostConstruct; import java.io.InputStream; Service public class PriceRuleEngine { Value(classpath:price-rules.json) private Resource ruleResource; private JsonNode rules; PostConstruct public void init() throws Exception { ObjectMapper mapper new ObjectMapper(); try (InputStream in ruleResource.getInputStream()) { rules mapper.readTree(in); } } public double calculate(String weather, double demandFactor, double supplyFactor) { JsonNode rule findRule(weather); if (rule null) { return rules.get(basePrice).asDouble(); } double base rules.get(basePrice).asDouble(); double demand rule.get(demandFactor).asDouble(); double supply rule.get(supplyFactor).asDouble(); // 示范公式价格 基础价 * 需求系数 / 供给系数 return base * demand * demandFactor / (supply * supplyFactor); } private JsonNode findRule(String weather) { for (JsonNode rule : rules.get(rules)) { if (weather.equals(rule.get(condition).asText())) { return rule; } } return null; } }最后提供一个价格查询接口验证效果。// 文件路径src/main/java/com/example/videodemo/controller/PriceController.java package com.example.videodemo.controller; import com.example.videodemo.common.Result; import com.example.videodemo.service.PriceRuleEngine; import com.example.videodemo.service.WeatherService; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; import javax.annotation.Resource; RestController RequestMapping(/price) public class PriceController { Resource private WeatherService weatherService; Resource private PriceRuleEngine priceRuleEngine; GetMapping(/query) public ResultDouble queryPrice(RequestParam String city, RequestParam(defaultValue 1.0) double demandFactor, RequestParam(defaultValue 1.0) double supplyFactor) { String weather weatherService.getWeatherByCity(city); double price priceRuleEngine.calculate(weather, demandFactor, supplyFactor); return Result.success(price); } }这样一个简化版动态定价链路就跑通了城市天气 - 匹配规则 - 结合供需系数计算价格。实际业务中供需系数可能来自订单系统、骑手/运力系统比这里复杂得多。4.8 运行验证启动VideoDemoApplication后可以通过接口测试分片上传curl -F fileIdtest123 -F chunkIndex0 -F totalChunks2 -F file/path/to/part0 http://localhost:8080/upload/chunk curl -F fileIdtest123 -F chunkIndex1 -F totalChunks2 -F file/path/to/part1 http://localhost:8080/upload/chunk合并分片curl http://localhost:8080/upload/merge?fileIdtest123fileNametest.mp4查询价格curl http://localhost:8080/price/query?citybeijing如果本地没有真实天气 API可以先在 WeatherService 中写死返回rain验证规则引擎是否按预期工作。5. 常见问题与排查思路问题现象常见原因解决思路分片上传后合并时报分片缺失Redis 分片状态与磁盘分片不一致检查 Redis 是否清空过 key确认分片目录是否被清理在前端记录已上传分片做补偿上传大文件时请求超时Nginx 或网关超时时间过短将超时时间调大或使用分片后限制单个分片大小转码任务堆积RabbitMQ 消息积压消费者能力不足或转码执行时间过长增加消费者实例或对转码任务做优先级分级限流不生效Sentinel 规则未加载或注解位置错误检查 Sentinel 控制台是否能看到资源确认规则已加载动态定价结果与预期不符规则条件未匹配或天气返回结构解析错误打印天气返回结果日志检查规则 JSON 格式MinIO 上传报 AccessDeniedBucket 未创建或密钥不正确先创建 Bucket确认 access-key 和 secret-key 无误排查这类问题时最重要的是把链路日志打通。比如分片上传链路要能看到每个分片的 Redis 记录、磁盘写入记录转码链路要能看到消息发送、消费、命令执行结果。日志缺失时很多问题只能靠猜。6. 最佳实践与工程建议6.1 分片大小与并发数分片大小不是越大越好也不是越小越好。常见经验值如下普通文件5MB 分片5-10 个并发。大视频文件10MB 分片10-20 个并发。弱网环境前端建议降低并发并支持暂停/恢复。分片过小会导致请求数量激增增加服务端压力分片过大又回到“大请求”的老问题。建议先小范围测试确定你当前网络和服务端性能下的最优值。6.2 上传接口的幂等性分片上传接口天然存在幂等需求。用户可能因为网络抖动重复提交同一个分片或者前端重试时重复上传。在设计时要保证同一个 fileId 下相同 chunkIndex 的重复上传不会导致分片损坏。合并接口要防止并发触发避免文件被同时合并多次。可以在 Redis 中使用布尔标记或 SETNX 进行幂等控制也可以在数据库表中建立唯一索引。6.3 削峰填谷午高峰两小时的流量压力可以通过以下方式缓解上传接口异步化只接收分片和记录状态合并动作延后处理。转码任务排队执行设置队列长度告警。开启弹性扩容对无状态服务按 CPU 和 QPS 指标自动扩缩容。对播放侧提前在低峰期做好热门内容预热。6.4 动态定价的灰度与监控动态定价直接影响用户感知必须谨慎上线先小流量灰度比如按城市、按用户分组逐步放开。记录每一次定价决策的完整上下文天气、供需系数、命中规则、最终价格。建立异常告警比如价格超过历史均值 2 倍时触发人工确认。价格不是越贵越好也不是越便宜越好。规则引擎的价值在于让策略可配置、可回放、可审计。6.5 安全与权限上传接口要注意对文件类型做白名单校验防止上传恶意文件。对文件内容做病毒扫描或至少做格式探测。控制单文件大小避免存储资源被恶意打满。对外暴露的定价接口要加权限校验不能在无认证情况下被刷。7. 结语与后续扩展回到标题里的三个关键点超长视频、午高峰两小时、恶劣天气单价。它们分别对应了存储与上传稳定性、高并发削峰能力、动态策略计算能力。这套代码示例虽然简化了很多细节但已经把核心链路打通了分片上传 - 状态记录 - 合并 - 异步转码 - 限流保护 - 天气数据接入 - 价格计算。你在实际项目中继续深入时可以优先补这几个方向将分片上传的 Redis 状态迁移到数据库保证更可靠的持久化。引入配置中心动态调整限流阈值和价格规则。完善监控大盘覆盖上传成功率、队列积压量、转码耗时、定价命中率等指标。如果视频量继续上涨把转码任务从单机执行迁移到大数据任务平台或容器编排平台。希望这篇文章能帮你把零散的需求变成可落地的技术方案。如果后续你在这个方向上有更具体的踩坑经验欢迎在评论区分享大家一起完善这套方案。
上一篇/下一篇内容由系统自动关联 返回资讯列表 →