mobile wallpaper 1mobile wallpaper 2mobile wallpaper 3mobile wallpaper 4mobile wallpaper 5mobile wallpaper 6mobile wallpaper 7mobile wallpaper 8mobile wallpaper 9mobile wallpaper 10mobile wallpaper 11mobile wallpaper 12
428 字
1 分钟
优化:将针对单一日志表的冷热数据分离类改造成通用类
2025-07-07

一、旧版方案#

点餐场景下:分析实现十万级用户日志的店铺推荐方案-CSDN博客


二、分析#

旧版方案该是将热数据存储到ES,由于我们新版方案考虑到MySQL和ES的查询差异,所以需要改为冷数据存储到ES。

此外,原定时任务的代码存在冗余问题。由于其操作的 和数据表是固定的,如果需要对另一个日志表实现相同的存储逻辑,就只能复制粘贴代码,这会导致大量的重复代码。

本次改进的核心目标就是解决这个问题,提升代码复用性。

如何将其改造为通用?思路其实很清晰:

  1. 将原代码中固定不变的部分参数化
  2. 引入泛型机制,使组件能够适配不同类型的数据实体。
  3. 对于 ES 操作部分,由于其模板化程度高,且不同日志类型可能需要不同的插入逻辑(如索引名、映射规则),我们将为每个日志表提供专属的重载函数以满足定制化需求。

文字描述可能不够直观,下面通过具体的来展示改进后的实现方式:


三、具体流程图#


四、代码实现#

@Slf4j
@RequiredArgsConstructor
@Component
public 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
作者
Yilena
发布于
2025-07-07
许可协议
CC BY 4.0

部分信息可能已经过时

相关文章 智能推荐
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代码实现,为海量日志分析与推荐系统落地提供了实践参考。

目录

封面
Sample Song
Sample Artist
封面
Sample Song
Sample Artist
0:00 / 0:00