news 2026/8/7 16:16:53

Spring Boot高并发场景下用户观看记录模块设计与实现

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Spring Boot高并发场景下用户观看记录模块设计与实现

最近在开发一个基于Spring Boot的在线视频平台项目时,遇到了一个非常典型的场景:如何优雅地处理一个业务模块(比如“用户观看记录”)从数据采集、处理、存储到前端展示的完整流程。这让我想起了一个有趣的比喻——就像记录“宁姆韦德普通的一天”,看似简单,实则涉及后端服务、数据库、缓存、消息队列乃至前端组件的协同工作。本文将围绕这个业务场景,拆解其技术实现,手把手带你构建一个高可用、可扩展的观看记录功能模块。

本文适合有一定Spring Boot和MyBatis基础的开发者,无论是想学习如何设计一个完整的业务闭环,还是希望优化现有项目的类似功能,都能从中获得启发。我们将从需求分析、表设计开始,逐步完成核心服务、异步处理、缓存策略的实现,并最终提供一个可复用的前端组件思路。

1. 业务背景与核心概念

在视频平台中,“观看记录”是一个基础但至关重要的功能。它不仅仅是记录用户看了什么,更关联着个性化推荐、内容热度计算、用户行为分析等多个下游业务。

1.1 核心价值与挑战

  • 用户体验:方便用户续播、查找历史。
  • 业务智能:为推荐系统提供原始数据。
  • 技术挑战
    • 高并发写入:热门视频同时有成千上万人观看。
    • 实时性要求:用户希望记录能即时同步到所有设备。
    • 数据一致性:记录进度需要准确,避免跳错时间点。
    • 存储成本:用户观看行为频繁,数据量增长快。

1.2 业务流程拆解一个完整的“记录”动作,可以分解为以下几个步骤:

  1. 事件触发:前端播放器每隔一段时间(如15秒)或暂停、退出时上报进度。
  2. 请求接收:后端API接收上报数据。
  3. 业务处理:清洗、验证数据,补充业务信息(如视频标题)。
  4. 数据持久化:将记录存入数据库。为了应对高并发,此处常引入异步和批处理。
  5. 缓存更新:更新用户最新的观看记录缓存,供快速查询。
  6. 下游通知:可选。通过消息队列通知推荐、统计等服务。

接下来,我们将从环境搭建开始,一步步实现这个流程。

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秒上报一次)策略。
  • 离线记录:在弱网环境下,可将记录暂存于浏览器的IndexedDBlocalStorage,待网络恢复后同步。
  • 进度同步:在多端(Web、App、TV)观看时,通过WebSocket或轮询及时同步最新进度,提供无缝体验。

通过以上步骤,我们完成了一个从需求分析到代码实现,再到优化扩展的“观看记录”功能模块。它不再是一个简单的INSERT语句,而是一个考虑了并发、性能、一致性的小型系统。在实际项目中,你需要根据业务规模和技术架构做出权衡和选择。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/7 16:16:18

ESP-IDF物联网开发框架:从核心架构到实战应用

1. ESP-IDF&#xff1a;物联网开发的“瑞士军刀” 如果你正在寻找一个能让你从零开始&#xff0c;快速构建稳定、功能丰富的物联网设备的开发框架&#xff0c;那么ESP-IDF绝对是你绕不开的核心工具。它不是什么遥不可及的“黑科技”&#xff0c;而是乐鑫官方为自家ESP32、ESP32…

作者头像 李华
网站建设 2026/8/7 16:15:23

RevokeMsgPatcher终极指南:Windows版微信/QQ/TIM防撤回完整教程

RevokeMsgPatcher终极指南&#xff1a;Windows版微信/QQ/TIM防撤回完整教程 【免费下载链接】RevokeMsgPatcher :trollface: A hex editor for WeChat/QQ/TIM - PC版微信/QQ/TIM防撤回补丁&#xff08;我已经看到了&#xff0c;撤回也没用了&#xff09; 项目地址: https://g…

作者头像 李华
网站建设 2026/8/7 16:15:14

Mermaid Live Editor:3分钟上手,免费创建专业图表的神奇工具

Mermaid Live Editor&#xff1a;3分钟上手&#xff0c;免费创建专业图表的神奇工具 【免费下载链接】mermaid-live-editor Edit, preview and share mermaid charts/diagrams. New implementation of the live editor. 项目地址: https://gitcode.com/GitHub_Trending/me/me…

作者头像 李华
网站建设 2026/8/7 16:15:11

HsMod:炉石传说玩家的终极效率工具箱

HsMod&#xff1a;炉石传说玩家的终极效率工具箱 【免费下载链接】HsMod Hearthstone Modification Based on BepInEx 项目地址: https://gitcode.com/GitHub_Trending/hs/HsMod HsMod 是一款基于 BepInEx 框架开发的炉石传说多功能插件&#xff0c;专为提升游戏体验和操…

作者头像 李华
网站建设 2026/8/7 16:14:58

STM32+FreeRTOS嵌入式开发实战:从环境搭建到多任务系统设计

1. 从零到一&#xff1a;为什么STM32FreeRTOS是嵌入式开发的“硬通货” 如果你刚开始接触单片机&#xff0c;或者已经玩过51、Arduino&#xff0c;现在想往更专业的嵌入式领域走&#xff0c;那么STM32和FreeRTOS这个组合&#xff0c;就是你绕不开的“硬通货”。它解决的核心问题…

作者头像 李华