428 字
1 分钟
优化:将针对单一日志表的冷热数据分离类改造成通用类
一、旧版方案
点餐场景下:分析实现十万级用户日志的店铺推荐方案-CSDN博客
二、分析
旧版方案该是将热数据存储到ES,由于我们新版方案考虑到MySQL和ES的查询差异,所以需要改为冷数据存储到ES。
此外,原定时任务的代码存在冗余问题。由于其操作的 和数据表是固定的,如果需要对另一个日志表实现相同的存储逻辑,就只能复制粘贴代码,这会导致大量的重复代码。
本次改进的核心目标就是解决这个问题,提升代码复用性。
如何将其改造为通用?思路其实很清晰:
- 将原代码中固定不变的部分参数化。
- 引入泛型机制,使组件能够适配不同类型的数据实体。
- 对于 ES 操作部分,由于其模板化程度高,且不同日志类型可能需要不同的插入逻辑(如索引名、映射规则),我们将为每个日志表提供专属的重载函数以满足定制化需求。
文字描述可能不够直观,下面通过具体的来展示改进后的实现方式:
三、具体流程图

四、代码实现
@Slf4j@RequiredArgsConstructor@Componentpublic class ESHotAndColdDataSeparationTask {
private final RestHighLevelClient restHighLevelClient; private final EnterStoreUserLogMapper enterStoreUserLogMapper; private final EnterStoreLogMapper enterStoreLogMapper; private final EnterStoreByCardLogMapper enterStoreByCardLogMapper; private final EnterStoreUserByCardLogMapper enterStoreUserByCardLogMapper; private final DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyyMMdd"); private final ThreadPoolExecutor threadPoolExecutor;
@XxlJob( value = "es-enterStoreUserLog-hotAndColdDataSeparation") public void processHotAndColdDataSeparation() { LocalDateTime start = LocalDateTime.now(); log.info("ES热数据表和冷数据表分离定时任务开始执行,执行时间:{}", start);
LocalDateTime twoMonthsAgo = LocalDateTime.now().minusMonths(2) .withHour(0).withMinute(0).withSecond(0);
List<CompletableFuture<Void>> tasks = List.of( // 处理 EnterStoreUserLog createTask(() -> dealLogs( enterStoreUserLogMapper, EnterStoreUserLog.class, EnterStoreUserLog::getId, EnterStoreUserLog::getCreateTime, twoMonthsAgo, this::processEnterStoreUserLogToEs ), new RuntimeException("处理EnterStoreUserLog失败"), TaskConstant.ENTER_STORE_USER_LOG,threadPoolExecutor),
// 处理 EnterStoreLog createTask(() -> dealLogs( enterStoreLogMapper, EnterStoreLog.class, EnterStoreLog::getId, EnterStoreLog::getCreateTime, twoMonthsAgo, this::processEnterStoreLogToEs ),new RuntimeException("处理EnterStoreLog失败"), TaskConstant.ENTER_STORE_LOG,threadPoolExecutor),
// 处理 EnterStoreUserByCardLog createTask(() -> dealLogs( enterStoreUserByCardLogMapper, EnterStoreUserByCardLog.class, EnterStoreUserByCardLog::getId, EnterStoreUserByCardLog::getCreateTime, twoMonthsAgo, this::processEnterStoreUserByCardLogToEs ),new RuntimeException("处理EnterStoreUserByCardLog失败"), TaskConstant.ENTER_STORE_USER_BY_CARD_LOG,threadPoolExecutor),
// 处理 EnterStoreByCardLog createTask(() -> dealLogs( enterStoreByCardLogMapper, EnterStoreByCardLog.class, EnterStoreByCardLog::getId, EnterStoreByCardLog::getCreateTime, twoMonthsAgo, this::processEnterStoreByCardLogToEs ),new RuntimeException("处理EnterStoreByCardLog失败"), TaskConstant.ENTER_STORE_BY_CARD_LOG,threadPoolExecutor) );
try{ // 等待所有任务完成 CompletableFuture.allOf(tasks.toArray(new CompletableFuture[0])).join(); } catch (Exception e) { throw new RuntimeException("多线程处理数据时发生异常", e); }
log.info("ES热数据表和冷数据表分离定时任务执行完毕,结束时间:{},花费时间:{}s", LocalDateTime.now(), Duration.between(start, LocalDateTime.now()).getSeconds()); }
/** * 通用日志处理方法 * @param mapper 实体对应的Mapper * @param clazz 实体类型 * @param idField 获取ID字段的函数 * @param createTimeField 获取创建时间字段的函数 * @param offsetDate 多久之前的数据 * @param esProcessor 插入es的重载函数 */ private <T> void dealLogs( BaseMapper<T> mapper, Class<T> clazz, SFunction<T, Long> idField, SFunction<T, LocalDateTime> createTimeField, LocalDateTime offsetDate, Function<List<T>, Boolean> esProcessor) {
// 获取初始游标(最早的一条记录) LambdaQueryWrapper<T> initWrapper = Wrappers.lambdaQuery(clazz) .lt(createTimeField, offsetDate) .orderByAsc(createTimeField) .orderByAsc(idField) .last("LIMIT 1");
// 查询 T first = mapper.selectOne(initWrapper);
// 如果为空则代表无数据,直接返回即可 if (first == null) { log.info("无需要迁移的{}数据", clazz.getSimpleName()); return; }
// 调用对应的函数获取id Long cursor = idField.apply(first);
// 开始查询 while (cursor != null) { // 构建查询条件 LambdaQueryWrapper<T> queryWrapper = Wrappers.lambdaQuery(clazz) .gt(idField, cursor) .lt(createTimeField, offsetDate) .orderByAsc(createTimeField) .orderByAsc(idField) .last("LIMIT 1000");
// 查询 List<T> batch = mapper.selectList(queryWrapper);
// 如果无数据则代表查询结束 if (batch.isEmpty()) break;
// 批量插入 boolean esSuccess = executeWithRetry(() -> esProcessor.apply(batch));
// 如果无异常抛出 if (esSuccess) { // 批量删除MySQL数据 List<Long> ids = batch.stream() .map(idField) .collect(Collectors.toList()); mapper.deleteBatchIds(ids); } else { // 如果删除失败就保留,进行下次循环 log.error("ES插入失败,保留MySQL数据: {}", clazz.getSimpleName()); }
// 更新游标 T last = batch.getLast(); cursor = idField.apply(last); } }
/** * 创建任务 * @param action 需要执行的任务 * @param exception 异常 * @param taskName 任务名称 * @param executor 线程池 * @return CompletableFuture */ private CompletableFuture<Void> createTask(Runnable action, RuntimeException exception, String taskName, Executor executor) { return CompletableFuture.runAsync(() -> { try { action.run(); } catch (Exception e) { log.error("{} 执行失败", taskName, e); if (exception != null) { throw exception; } } }, executor); }
/** * 重试机制 */ private boolean executeWithRetry(Supplier<Boolean> operation) { try { // 返回插入数据结果 return operation.get(); } catch (Exception e) { try { // 短暂休眠防止网络波动导致插入失败 Thread.sleep(200); // 二次重试 return operation.get(); } catch (Exception ex) { // 二次重试再失败则记录日志,进行下次循环,因为如果插入失败的话就不会进行MySQL的删除操作 log.error("操作二次失败", ex); return false; } } }
/** * EnterStoreUserLog的ES处理 */ private Boolean processEnterStoreUserLogToEs(List<EnterStoreUserLog> batch) { try { // 构建批量插入请求 BulkRequest bulkRequest = new BulkRequest();
// 批量插入 for (EnterStoreUserLog log : batch) { String date = log.getCreateTime().toLocalDate().format(formatter); String esId = log.getId().toString();
IndexRequest request = new IndexRequest(ESIndexConstant.ENTER_STORE_USER_LOG + "_" + date) .id(esId) .source(JSONUtil.toJsonStr(log), XContentType.JSON);
bulkRequest.add(request); } restHighLevelClient.bulk(bulkRequest, RequestOptions.DEFAULT); return true; } catch (IOException e) { throw new RuntimeException("ES插入失败", e); } }
/** * EnterStoreLog的ES处理 */ private Boolean processEnterStoreLogToEs(List<EnterStoreLog> batch) { try { BulkRequest bulkRequest = new BulkRequest();
for (EnterStoreLog log : batch) { String date = log.getCreateTime().toLocalDate().format(formatter); String esId = log.getId().toString();
IndexRequest request = new IndexRequest(ESIndexConstant.ENTER_STORE_LOG + "_" + date) .id(esId) .source(JSONUtil.toJsonStr(log), XContentType.JSON);
bulkRequest.add(request); } restHighLevelClient.bulk(bulkRequest, RequestOptions.DEFAULT); return true; } catch (IOException e) { throw new RuntimeException("ES插入失败", e); } }
/** * EnterStoreUserByCardLog的ES处理 */ private Boolean processEnterStoreUserByCardLogToEs(List<EnterStoreUserByCardLog> batch) { try { // 构建批量插入请求 BulkRequest bulkRequest = new BulkRequest();
// 批量插入 for (EnterStoreUserByCardLog log : batch) { String date = log.getCreateTime().toLocalDate().format(formatter); String esId = log.getId().toString();
IndexRequest request = new IndexRequest(ESIndexConstant.ENTER_STORE_USER_BY_CARD_LOG + "_" + date) .id(esId) .source(JSONUtil.toJsonStr(log), XContentType.JSON);
bulkRequest.add(request); } restHighLevelClient.bulk(bulkRequest, RequestOptions.DEFAULT); return true; } catch (IOException e) { throw new RuntimeException("ES插入失败", e); } }
/** * EnterStoreByCardLog的ES处理 */ private Boolean processEnterStoreByCardLogToEs(List<EnterStoreByCardLog> batch) { try { BulkRequest bulkRequest = new BulkRequest();
for (EnterStoreByCardLog log : batch) { String date = log.getCreateTime().toLocalDate().format(formatter); String esId = log.getId().toString();
IndexRequest request = new IndexRequest(ESIndexConstant.ENTER_STORE_BY_CARD_LOG + "_" + date) .id(esId) .source(JSONUtil.toJsonStr(log), XContentType.JSON);
bulkRequest.add(request); } restHighLevelClient.bulk(bulkRequest, RequestOptions.DEFAULT); return true; } catch (IOException e) { throw new RuntimeException("ES插入失败", e); } }}如果你有更好的方案,请在评论区告诉我!
码文不易,点个赞再走吧
原文链接: 优化:将针对单一日志表的冷热数据分离类改造成通用类 作者: Yilena
分享
如果这篇文章对你有帮助,欢迎分享给更多人!
优化:将针对单一日志表的冷热数据分离类改造成通用类
https://blog.csdn.net/2401_88959292/article/details/148619523?spm=1001.2014.3001.5501 部分信息可能已经过时
相关文章 智能推荐
1
优化日志分析店铺推荐方案:用户范围的精确度以及ES与MySQL的查询效率差异
业务拆解 本文针对点餐场景下的店铺推荐方案进行了深度优化。首先,通过Redis记录用户月度登录天数,精准筛选出热点用户,解决了旧版方案中推荐用户范围不精确的问题。其次,对比了Elasticsearch与MySQL在十万级数据量下的查询效率,将热数据迁移至MySQL并建立索引以提升查询性能,同时保留ES作为冷数据备份。文章详细展示了优化后的业务流程图及Java代码实现,包括基于Lua脚本的Redis原子操作、多线程并行处理用户日志分析以及加权评分推荐算法,为高并发场景下的日志分析与个性化推荐提供了高效的工程实践。
2
布隆过滤器因内存上限无法处理超量数据的优化方案
业务拆解 本文针对布隆过滤器因数据量增加导致内存上限不足的问题,深入分析并提出了三种优化方案。首先是分片布隆过滤器,通过哈希取模将数据分散,简单高效但受限于单机内存且易引发GC问题;其次是分布式布隆过滤器,借助Redis摆脱单机限制,并结合一致性哈希实现分片以优化缓存空间;最后是可扩展布隆过滤器,通过动态增加层数和容量来应对数据增长,但查询效率会随层数增加而降低。综合对比,推荐单机项目使用分片方案,分布式架构采用Redis分布式分片方案,以有效解决超量数据处理难题。
3
116秒→6秒:Redis管道+批处理优化用户好友关系校验的方案
业务拆解 本文针对社交平台中用户好友关系数据一致性问题,提出了一种高效的定时任务解决方案。通过分析初版方案的性能瓶颈(单线程串行处理导致116秒耗时),逐步优化为多线程并行处理(70秒)和最终版批量预加载策略(6秒)。终版方案的核心改进包括:1)预加载所有关注关系并建立内存映射;2)批量处理好友数据更新;3)使用Redis管道技术减少网络请求。最终将请求次数从35万次降至常数级,同时提供了完整的Java实现代码,包含分片处理、批量数据库操作和Redis管道更新等关键优化技术。
4
360s→15s:近25倍的性能优化重构十万Excel券码批量导入方案
业务拆解 本文针对十万级Excel券码批量导入场景,将原逐行解析导致的20万次网络IO,通过Redis管道与MQ批处理优化至182次。方案兼顾了宕机恢复与库存扣减,最终将耗时从360秒大幅缩减至15秒左右。
5
点餐场景下:分析实现十万级用户日志的店铺推荐方案
业务拆解 本文针对点餐小程序十万级用户进店日志场景,设计并实现了一套高效的个性化店铺推荐方案。通过冷热数据分离策略,利用Elasticsearch存储近30天热数据,结合MySQL与Redis优化查询性能。文章详细阐述了基于访问频次、分区及分类偏好的加权评分排序算法,并规划了每日定时执行的离线分析任务以平衡系统开销。最后,提供了完整的实体类定义、ES数据迁移定时任务及多线程并发日志分析的Java代码实现,为海量日志分析与推荐系统落地提供了实践参考。










