最近在开发一个基于Spring Boot的在线视频平台项目时,遇到了一个非常典型的场景:如何优雅地处理一个业务模块(比如“用户观看记录”)从数据采集、处理、存储到前端展示的完整流程。这让我想起了一个有趣的比喻——就像记录“宁姆韦德普通的一天”,看似简单,实则涉及后端服务、数据库、缓存、消息队列乃至前端组件的协同工作。本文将围绕这个业务场景,拆解其技术实现,手把手带你构建一个高可用、可扩展的观看记录功能模块。
本文适合有一定Spring Boot和MyBatis基础的开发者,无论是想学习如何设计一个完整的业务闭环,还是希望优化现有项目的类似功能,都能从中获得启发。我们将从需求分析、表设计开始,逐步完成核心服务、异步处理、缓存策略的实现,并最终提供一个可复用的前端组件思路。
1. 业务背景与核心概念
在视频平台中,“观看记录”是一个基础但至关重要的功能。它不仅仅是记录用户看了什么,更关联着个性化推荐、内容热度计算、用户行为分析等多个下游业务。
1.1 核心价值与挑战
- 用户体验:方便用户续播、查找历史。
- 业务智能:为推荐系统提供原始数据。
- 技术挑战:
- 高并发写入:热门视频同时有成千上万人观看。
- 实时性要求:用户希望记录能即时同步到所有设备。
- 数据一致性:记录进度需要准确,避免跳错时间点。
- 存储成本:用户观看行为频繁,数据量增长快。
1.2 业务流程拆解一个完整的“记录”动作,可以分解为以下几个步骤:
- 事件触发:前端播放器每隔一段时间(如15秒)或暂停、退出时上报进度。
- 请求接收:后端API接收上报数据。
- 业务处理:清洗、验证数据,补充业务信息(如视频标题)。
- 数据持久化:将记录存入数据库。为了应对高并发,此处常引入异步和批处理。
- 缓存更新:更新用户最新的观看记录缓存,供快速查询。
- 下游通知:可选。通过消息队列通知推荐、统计等服务。
接下来,我们将从环境搭建开始,一步步实现这个流程。
2. 环境准备与项目结构
我们使用当前主流的Java技术栈进行演示。
2.1 基础环境
- JDK: 17 或以上 (推荐17,长期支持版本)
- Maven: 3.6+
- IDE: IntelliJ IDEA 或 VS Code
- 数据库: MySQL 8.0
- 缓存: Redis 6.x
2.2 项目初始化与依赖使用 Spring Initializr 创建一个Spring Boot项目,选择以下依赖:
- Spring Web: 提供RESTful API支持。
- Spring Data Redis: 操作Redis缓存。
- MyBatis Framework: 数据库ORM框架。
- MySQL Driver: 连接MySQL数据库。
- Lombok: 简化Java Bean代码。
生成的pom.xml关键依赖如下:
<dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-redis</artifactId> </dependency> <dependency> <groupId>org.mybatis.spring.boot</groupId> <artifactId>mybatis-spring-boot-starter</artifactId> <version>3.0.3</version> <!-- 请使用最新稳定版 --> </dependency> <dependency> <groupId>com.mysql</groupId> <artifactId>mysql-connector-j</artifactId> <scope>runtime</scope> </dependency> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> <optional>true</optional> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-test</artifactId> <scope>test</scope> </dependency> </dependencies>2.3 配置文件配置application.yml,设置数据源、Redis和MyBatis。
server: port: 8080 spring: datasource: url: jdbc:mysql://localhost:3306/video_platform?useUnicode=true&characterEncoding=utf8&serverTimezone=Asia/Shanghai username: your_username password: your_password driver-class-name: com.mysql.cj.jdbc.Driver redis: host: localhost port: 6379 password: '' # 如果有密码则填写 database: 0 lettuce: pool: max-active: 8 max-wait: -1ms max-idle: 8 min-idle: 0 mybatis: mapper-locations: classpath:mapper/*.xml configuration: map-underscore-to-camel-case: true # 开启驼峰命名自动转换2.4 项目结构预览
src/main/java/com/example/videoplatform/ ├── VideoPlatformApplication.java ├── config/ # 配置类 ├── controller/ # 控制层,接收API请求 ├── service/ # 业务逻辑层 │ ├── impl/ ├── mapper/ # MyBatis Mapper接口 ├── entity/ # 实体类,对应数据库表 ├── dto/ # 数据传输对象 ├── vo/ # 视图对象,用于接口返回 └── async/ # 异步处理组件3. 数据库设计与实体建模
观看记录的核心在于表结构设计,需平衡查询效率与存储空间。
3.1 表结构设计 (SQL)
-- 用户观看记录表 CREATE TABLE `user_watch_history` ( `id` bigint(20) NOT NULL AUTO_INCREMENT COMMENT '主键ID', `user_id` bigint(20) NOT NULL COMMENT '用户ID', `video_id` bigint(20) NOT NULL COMMENT '视频ID', `watch_progress` int(11) NOT NULL DEFAULT '0' COMMENT '观看进度(秒)', `video_duration` int(11) NOT NULL COMMENT '视频总时长(秒)', `latest_watch_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '最近观看时间', `created_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '记录创建时间', `updated_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '记录更新时间', `is_deleted` tinyint(1) NOT NULL DEFAULT '0' COMMENT '逻辑删除标志', PRIMARY KEY (`id`), -- 唯一索引,一个用户对同一个视频只保留一条最新记录 UNIQUE KEY `uk_user_video` (`user_id`,`video_id`), -- 用于查询用户的历史记录列表 KEY `idx_user_time` (`user_id`,`latest_watch_time`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='用户观看记录表';设计要点:
uk_user_video唯一索引:确保一个用户对一个视频只有一条记录,更新时使用ON DUPLICATE KEY UPDATE或先查后改,避免数据膨胀。idx_user_time索引:优化按用户和时间倒序查询列表的性能。is_deleted:逻辑删除标志,避免物理删除。
3.2 实体类 (Entity)对应上述表结构,创建Java实体类。
// 文件路径:src/main/java/com/example/videoplatform/entity/UserWatchHistory.java package com.example.videoplatform.entity; import lombok.Data; import java.time.LocalDateTime; @Data public class UserWatchHistory { private Long id; private Long userId; private Long videoId; private Integer watchProgress; // 单位:秒 private Integer videoDuration; // 单位:秒 private LocalDateTime latestWatchTime; private LocalDateTime createdTime; private LocalDateTime updatedTime; private Boolean isDeleted; }3.3 数据传输对象 (DTO) 和视图对象 (VO)
- DTO (WatchProgressDTO):用于接收前端上报的进度数据。
@Data public class WatchProgressDTO { @NotNull(message = "视频ID不能为空") private Long videoId; @Min(value = 0, message = "进度不能小于0") private Integer progress; // 当前播放进度(秒) @Min(value = 1, message = "时长必须大于0") private Integer duration; // 视频总时长(秒) } - VO (WatchHistoryVO):用于返回给前端的观看记录信息,通常会关联视频信息。
@Data public class WatchHistoryVO { private Long videoId; private String videoTitle; private String coverUrl; private Integer watchProgress; private Integer videoDuration; private String latestWatchTime; // 格式化后的时间字符串 // 可以计算一个进度百分比,方便前端显示 public String getProgressPercentage() { if (videoDuration == null || videoDuration == 0) return "0%"; double percentage = (watchProgress.doubleValue() / videoDuration) * 100; return String.format("%.1f%%", Math.min(percentage, 100)); } }
4. 核心业务逻辑实现
我们将采用“异步处理 + 缓存”的策略来应对高并发写入和实时查询。
4.1 Mapper 接口与 XML首先定义数据访问层。
// 文件路径:src/main/java/com/example/videoplatform/mapper/UserWatchHistoryMapper.java @Mapper public interface UserWatchHistoryMapper { // 插入或更新记录(使用ON DUPLICATE KEY UPDATE) int upsert(UserWatchHistory history); // 查询用户最新的N条观看记录 List<UserWatchHistory> selectByUserId(@Param("userId") Long userId, @Param("limit") Integer limit); // 逻辑删除某条记录 int logicDelete(@Param("id") Long id, @Param("userId") Long userId); }对应的UserWatchHistoryMapper.xml:
<!-- 文件路径:src/main/resources/mapper/UserWatchHistoryMapper.xml --> <mapper namespace="com.example.videoplatform.mapper.UserWatchHistoryMapper"> <insert id="upsert" parameterType="UserWatchHistory"> INSERT INTO user_watch_history (user_id, video_id, watch_progress, video_duration, latest_watch_time) VALUES (#{userId}, #{videoId}, #{watchProgress}, #{videoDuration}, NOW()) ON DUPLICATE KEY UPDATE watch_progress = VALUES(watch_progress), video_duration = VALUES(video_duration), latest_watch_time = NOW(), updated_time = NOW() </insert> <select id="selectByUserId" resultType="UserWatchHistory"> SELECT * FROM user_watch_history WHERE user_id = #{userId} AND is_deleted = 0 ORDER BY latest_watch_time DESC LIMIT #{limit} </select> <update id="logicDelete"> UPDATE user_watch_history SET is_deleted = 1, updated_time = NOW() WHERE id = #{id} AND user_id = #{userId} </update> </mapper>4.2 服务层实现 (Service)服务层负责核心业务逻辑,这里我们引入异步处理。
// 文件路径:src/main/java/com/example/videoplatform/service/WatchHistoryService.java public interface WatchHistoryService { void recordWatchProgress(Long userId, WatchProgressDTO dto); List<WatchHistoryVO> getWatchHistory(Long userId, Integer limit); boolean deleteHistory(Long userId, Long recordId); }// 文件路径:src/main/java/com/example/videoplatform/service/impl/WatchHistoryServiceImpl.java @Service @Slf4j public class WatchHistoryServiceImpl implements WatchHistoryService { @Autowired private UserWatchHistoryMapper historyMapper; @Autowired private RedisTemplate<String, Object> redisTemplate; @Autowired private AsyncTaskExecutor asyncTaskExecutor; // 自定义的异步执行器 private static final String WATCH_HISTORY_KEY_PREFIX = "wh:uid:"; @Override public void recordWatchProgress(Long userId, WatchProgressDTO dto) { // 1. 参数校验 (略) // 2. 构造实体 UserWatchHistory history = new UserWatchHistory(); history.setUserId(userId); history.setVideoId(dto.getVideoId()); history.setWatchProgress(dto.getProgress()); history.setVideoDuration(dto.getDuration()); // 3. 异步执行数据库持久化 asyncTaskExecutor.execute(() -> { try { int rows = historyMapper.upsert(history); log.debug("观看记录持久化成功,userId:{}, videoId:{}, affected rows:{}", userId, dto.getVideoId(), rows); } catch (Exception e) { log.error("观看记录持久化失败,userId:{}, videoId:{}", userId, dto.getVideoId(), e); // 此处可加入降级策略,如存入本地队列重试或记录日志 } }); // 4. 同步更新Redis缓存 (保证实时性) String cacheKey = WATCH_HISTORY_KEY_PREFIX + userId; WatchHistoryVO cacheVO = new WatchHistoryVO(); // 这里需要从其他服务或数据库获取视频详情,简化演示 cacheVO.setVideoId(dto.getVideoId()); cacheVO.setWatchProgress(dto.getProgress()); cacheVO.setVideoDuration(dto.getDuration()); cacheVO.setLatestWatchTime(LocalDateTime.now().toString()); // 使用Hash结构存储,field为videoId redisTemplate.opsForHash().put(cacheKey, dto.getVideoId().toString(), cacheVO); // 设置缓存过期时间,例如7天 redisTemplate.expire(cacheKey, 7, TimeUnit.DAYS); } @Override public List<WatchHistoryVO> getWatchHistory(Long userId, Integer limit) { List<WatchHistoryVO> result = new ArrayList<>(); String cacheKey = WATCH_HISTORY_KEY_PREFIX + userId; // 1. 先查缓存 Map<Object, Object> cacheMap = redisTemplate.opsForHash().entries(cacheKey); if (cacheMap != null && !cacheMap.isEmpty()) { // 缓存存在,转换并排序 cacheMap.values().forEach(obj -> result.add((WatchHistoryVO) obj)); result.sort((a, b) -> b.getLatestWatchTime().compareTo(a.getLatestWatchTime())); if (limit != null && result.size() > limit) { return result.subList(0, limit); } return result; } // 2. 缓存不存在,查数据库 List<UserWatchHistory> dbList = historyMapper.selectByUserId(userId, limit != null ? limit : 50); if (dbList.isEmpty()) { return result; } // 3. 转换并填充视频详情 (此处简化,实际需调用视频服务) for (UserWatchHistory history : dbList) { WatchHistoryVO vo = convertToVO(history); // 假设的转换方法 result.add(vo); // 4. 异步回写缓存 redisTemplate.opsForHash().put(cacheKey, history.getVideoId().toString(), vo); } redisTemplate.expire(cacheKey, 7, TimeUnit.DAYS); return result; } // convertToVO 等方法省略... }4.3 异步执行器配置为了避免数据库写入阻塞主线程,我们配置一个专用的线程池。
// 文件路径:src/main/java/com/example/videoplatform/config/AsyncConfig.java @Configuration @EnableAsync public class AsyncConfig { @Bean("asyncTaskExecutor") public TaskExecutor asyncTaskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); // 核心线程数 executor.setCorePoolSize(5); // 最大线程数 executor.setMaxPoolSize(20); // 队列容量 executor.setQueueCapacity(1000); // 线程名前缀 executor.setThreadNamePrefix("WatchHistory-Async-"); // 拒绝策略:由调用线程直接执行 executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); executor.initialize(); return executor; } }4.4 控制层 (Controller)提供对外的REST API。
// 文件路径:src/main/java/com/example/videoplatform/controller/WatchHistoryController.java @RestController @RequestMapping("/api/watch-history") @Slf4j public class WatchHistoryController { @Autowired private WatchHistoryService watchHistoryService; @PostMapping("/record") public ResponseEntity<Void> recordProgress(@RequestBody @Valid WatchProgressDTO dto, @RequestHeader("X-User-Id") Long userId) { // 实际项目中,userId应从Token或Session中获取,此处简化 if (userId == null || userId <= 0) { return ResponseEntity.status(HttpStatus.UNAUTHORIZED).build(); } try { watchHistoryService.recordWatchProgress(userId, dto); return ResponseEntity.ok().build(); } catch (Exception e) { log.error("记录观看进度失败,userId:{}, dto:{}", userId, dto, e); return ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).build(); } } @GetMapping("/list") public ResponseEntity<List<WatchHistoryVO>> getHistory(@RequestParam(defaultValue = "20") Integer limit, @RequestHeader("X-User-Id") Long userId) { if (userId == null || userId <= 0) { return ResponseEntity.status(HttpStatus.UNAUTHORIZED).build(); } List<WatchHistoryVO> history = watchHistoryService.getWatchHistory(userId, limit); return ResponseEntity.ok(history); } @DeleteMapping("/{recordId}") public ResponseEntity<Void> deleteHistory(@PathVariable Long recordId, @RequestHeader("X-User-Id") Long userId) { boolean success = watchHistoryService.deleteHistory(userId, recordId); return success ? ResponseEntity.ok().build() : ResponseEntity.notFound().build(); } }5. 前端交互模拟与测试
后端完成后,我们需要验证API。这里使用curl命令和单元测试进行模拟。
5.1 上报观看进度 (模拟请求)
curl -X POST 'http://localhost:8080/api/watch-history/record' \ -H 'Content-Type: application/json' \ -H 'X-User-Id: 123' \ -d '{ "videoId": 1001, "progress": 125, "duration": 600 }'5.2 查询观看记录
curl -X GET 'http://localhost:8080/api/watch-history/list?limit=10' \ -H 'X-User-Id: 123'5.3 服务层单元测试示例
// 文件路径:src/test/java/com/example/videoplatform/service/WatchHistoryServiceTest.java @SpringBootTest @Slf4j class WatchHistoryServiceTest { @Autowired private WatchHistoryService watchHistoryService; @Test void testRecordAndGetHistory() { Long userId = 999L; WatchProgressDTO dto = new WatchProgressDTO(); dto.setVideoId(2001L); dto.setProgress(30); dto.setDuration(180); // 测试记录 watchHistoryService.recordWatchProgress(userId, dto); // 等待异步任务执行(测试环境可简单等待) try { Thread.sleep(1000); } catch (InterruptedException e) { } // 测试查询 List<WatchHistoryVO> history = watchHistoryService.getWatchHistory(userId, 10); Assertions.assertNotNull(history); Assertions.assertFalse(history.isEmpty()); Assertions.assertEquals(dto.getVideoId(), history.get(0).getVideoId()); log.info("测试通过,查询到记录:{}", history.get(0).getProgressPercentage()); } }6. 常见问题与排查思路
在实际开发和运维中,你可能会遇到以下问题:
| 问题现象 | 可能原因 | 排查思路与解决方案 |
|---|---|---|
| 记录上报成功,但查询不到或进度未更新 | 1. 异步任务执行失败。 2. Redis缓存未正确更新或已过期。 3. 数据库唯一键冲突导致更新失败。 | 1. 查看应用日志,搜索“观看记录持久化失败”。 2. 使用 redis-cli检查对应Key是否存在:HGETALL wh:uid:123。3. 检查数据库 user_watch_history表,确认数据是否存在及进度是否正确。 |
| 接口响应缓慢,尤其是记录上报接口 | 1. 数据库写入慢(如未建索引、锁表)。 2. Redis连接池耗尽或网络延迟高。 3. 异步线程池队列满,触发拒绝策略。 | 1. 使用EXPLAIN分析upsert语句。2. 监控Redis连接数和响应时间。 3. 调整异步线程池配置( CorePoolSize,QueueCapacity),或监控线程池状态。 |
| 缓存与数据库数据不一致 | 1. 缓存更新成功但数据库更新失败。 2. 缓存过期后,从数据库回写时数据已变。 | 1.保证最终一致性:异步任务失败后应有重试机制(如存入死信队列)。 2.使用较短的缓存过期时间(如30分钟),并考虑在更新数据库后主动刷新缓存。 |
| 高并发下,数据库压力大 | 即使异步,瞬时写入量也可能很大。 | 1.引入消息队列(如Kafka/RocketMQ),将记录先发往队列,由消费者批量写入数据库。 2.合并写入:在内存中暂存一段时间内的进度,合并为一次更新。 |
| 用户量巨大,Redis内存占用高 | 每个用户的记录都缓存。 | 1.限制缓存数量:每个用户只缓存最新的N条(如50条)。 2.使用更紧凑的数据结构:例如只缓存 videoId:progress的映射,其他信息懒加载。3.设置合理的过期策略。 |
7. 最佳实践与进阶优化
实现基础功能后,我们可以从性能、可靠性和可扩展性方面进行优化。
7.1 性能优化
- 数据库层面:
- 对
user_id,video_id,latest_watch_time建立联合索引,优化查询。 - 定期归档或清理很久之前(如一年前)的观看记录,可以迁移到历史表或冷存储。
- 对
- 缓存层面:
- 使用Redis Pipeline批量操作缓存,减少网络往返。
- 考虑使用Redis Sorted Set来存储用户观看记录,
score设置为观看时间戳,天然支持按时间排序,且可以方便地按范围查询和限制数量。
7.2 可靠性保障
- 异步任务可靠性:
- 将异步任务提交到持久化消息队列(如RocketMQ),确保即使应用重启,任务也不会丢失。
- 实现消费者端的幂等性处理,防止因重试导致的数据重复更新。
- 降级与熔断:
- 当Redis不可用时,应能降级为直接查询数据库,避免核心功能不可用。
- 使用 Resilience4j 或 Sentinel 对数据库调用进行熔断保护。
7.3 架构扩展
- 分库分表:当用户量达到千万甚至亿级,单表性能成为瓶颈。可按
user_id进行分片。 - 读写分离:将读请求(查询历史记录)路由到从库,减轻主库压力。
- 引入Elasticsearch:如果需要支持复杂的搜索(如按视频标题搜索观看记录),可以将记录同步到ES中。
7.4 前端优化建议
- 上报节流:避免每秒上报多次,可以使用防抖(暂停时上报)或节流(每15秒上报一次)策略。
- 离线记录:在弱网环境下,可将记录暂存于浏览器的
IndexedDB或localStorage,待网络恢复后同步。 - 进度同步:在多端(Web、App、TV)观看时,通过WebSocket或轮询及时同步最新进度,提供无缝体验。
通过以上步骤,我们完成了一个从需求分析到代码实现,再到优化扩展的“观看记录”功能模块。它不再是一个简单的INSERT语句,而是一个考虑了并发、性能、一致性的小型系统。在实际项目中,你需要根据业务规模和技术架构做出权衡和选择。