某电商企业用Java Reducer合并千万条订单数据减少内存压力提升处理速度实战案例
先说说这家电商公司到底遇到了啥麻烦。
一、那个让人头疼的订单汇总
小陈是某中型电商平台的数据工程师,今年双11刚过,公司的订单量直接爆了——单日订单突破千万条,日均累计数据更是达到了数千万条。
公司有一个报表系统,需要把分散在多个数据源的订单数据合并,算出每个用户的累计消费金额、购买次数、平均客单价等核心指标。
原来的方案是:
// 旧方案:直接把所有数据加载到内存
List<Order> allOrders = orderRepository.findAll(); // 几千万条
Map<String, OrderSummary> summaryMap = new HashMap<>();
for (Order order : allOrders) {
String userId = order.getUserId();
OrderSummary summary = summaryMap.getOrDefault(userId, new OrderSummary(userId));
summary.increaseSum(order.getAmount());
summary.increaseCount();
}
听起来很简单对吧?但问题来了——当数据量达到千万级别时,这个方案直接让服务器内存爆了。
内存消耗有多夸张?
我们来算笔账:
假设每条订单对象约 500 字节
1000万条订单 = 500MB 纯数据
加上 HashMap 的节点开销、字符串对象、包装类...
实际内存消耗轻松突破 3GB
服务器总共才 8GB 内存,光一个报表任务就把内存吃光了,其他业务全卡住。
小陈的老板在晨会上说:”这个月KPI完不成,大家年终奖金打水漂。”
二、Reducer 思路是如何诞生的
小陈开始研究解决方案。他在GitHub上翻了一堆资料,发现Hadoop MapReduce的编程模型有一个核心思想让他眼前一亮——分而治之。
想象一下,你要数一堆硬币有多少钱:
传统做法:把所有硬币倒在桌上,一个一个数 Reducer思路:先把硬币分成10小堆,每堆单独数完金额,最后把10个小计加在一起
这样你不需要一次性看到所有硬币,内存压力就小多了。
小陈决定用Java实现一个类似的思想,但不是用Hadoop那么重的框架,而是用原生的Java多线程 + 分块处理。
三、核心代码实现
第一步:定义数据模型
/**
* 订单数据模型
*/
@Data
@NoArgsConstructor
@AllArgsConstructor
public class Order {
private String orderId; // 订单编号
private String userId; // 用户ID
private BigDecimal amount; // 订单金额
private Integer itemCount; // 商品数量
private LocalDateTime createTime;
private String status; // 订单状态
}
/**
* 用户订单汇总结果
*/
@Data
@Builder
public class OrderSummary {
private String userId;
private long totalOrders; // 总订单数
private BigDecimal totalAmount; // 总金额
private BigDecimal avgAmount; // 平均客单价
private long totalItems; // 总商品数
/**
* 增量更新——这是Reducer的核心操作
* 每次只处理一小批数据,合并结果
*/
public void merge(OrderSummary other) {
this.totalOrders += other.getTotalOrders();
this.totalAmount = this.totalAmount.add(other.getTotalAmount());
this.totalItems += other.getTotalItems();
}
}
第二步:分批读取数据,避免一次性加载
/**
* 分页查询服务
* 关键:每次只从数据库取 10000 条,绝不贪婪
*/
@Service
public class OrderQueryService {
@Autowired
private OrderRepository orderRepository;
/**
* 分页分批获取订单
* @param pageSize 每批大小,10000是经验值
* @return 流式数据,用尽即弃
*/
public Stream<Order> batchQueryOrders(int pageSize) {
int pageNum = 0;
List<Order> batch;
// 循环分页查询,用流的方式返回
return Stream.generate(() -> {
batch = orderRepository.findByPage(pageNum, pageSize);
pageNum++;
return batch;
})
.takeWhile(batch -> !batch.isEmpty())
.flatMap(List::stream);
}
}
第三步:实现 Reducer 合并逻辑
这是整个方案最核心的部分。
/**
* 订单数据 Reducer
*
* 核心思想:
* 1. 数据分块输入(多个Map输出)
* 2. 每个Reducer只维护自己负责的Key范围
* 3. 相同Key的数据汇聚到同一个Reducer
* 4. 最终结果合并
*/
@Component
public class OrderReducer {
private static final Logger log = LoggerFactory.getLogger(OrderReducer.class);
/**
* 用户ID分片策略:按userId哈希值取模
* 假设有8个Reducer,每个负责1/8的用户
*/
private final int reducerCount = 8;
/**
* 每个Reducer的内存缓冲区大小(条)
* 超过这个数就触发局部聚合
*/
private static final int BUFFER_SIZE = 50000;
/**
* 执行Reducer合并
* @param orders 订单数据流(已分块)
* @return 汇总结果
*/
public Map<String, OrderSummary> reduce(Stream<Order> orders) {
// ========== 阶段1:按Reducer编号分发 ==========
// 每个Reducer维护自己负责的用户汇总
List<ConcurrentHashMap<String, OrderSummary>> reducers =
IntStream.range(0, reducerCount)
.mapToObj(i -> new ConcurrentHashMap<>())
.collect(Collectors.toList());
// ========== 阶段2:流式处理 + 局部聚合 ==========
orders.forEach(order -> {
int reducerIndex = getReducerIndex(order.getUserId());
ConcurrentHashMap<String, OrderSummary> reducer = reducers.get(reducerIndex);
// 局部聚合:相同userId的数据合并到同一个map entry
reducer.merge(order.getUserId(),
buildPartialSummary(order),
OrderSummary::merge);
// 防内存溢出:当缓冲区过大时,主动触发压缩
if (reducer.size() > BUFFER_SIZE) {
compressReducer(reducer, reducerIndex, reducers);
}
});
// ========== 阶段3:跨Reducer合并 ==========
// 把多个Reducer的结果最终合并
return mergeReducers(reducers);
}
/**
* 根据userId确定归属哪个Reducer
* 简单哈希取模,确保相同用户一定落在同一个Reducer
*/
private int getReducerIndex(String userId) {
return Math.abs(userId.hashCode()) % reducerCount;
}
/**
* 构建局部汇总(增量式)
*/
private OrderSummary buildPartialSummary(Order order) {
return OrderSummary.builder()
.userId(order.getUserId())
.totalOrders(1)
.totalAmount(order.getAmount())
.totalItems(order.getItemCount())
.build();
}
/**
* 压缩Reducer:将小用户数据沉淀,释放内存
* 这是减少内存压力的关键技巧
*/
private void compressReducer(
ConcurrentHashMap<String, OrderSummary> currentReducer,
int reducerIndex,
List<ConcurrentHashMap<String, OrderSummary>> allReducers) {
// 找出累计金额最小的50%用户,先行输出
List<String> toPersist = currentReducer.entrySet().stream()
.filter(e -> e.getValue().getTotalAmount().compareTo(BigDecimal.valueOf(100)) < 0)
.limit(currentReducer.size() / 2)
.map(Map.Entry::getKey)
.collect(Collectors.toList());
// 持久化到临时文件(防OOM的安全阀)
if (!toPersist.isEmpty()) {
persistToTempFile(reducerIndex, toPersist, currentReducer);
toPersist.forEach(currentReducer::remove);
log.info("Reducer-{} 完成内存压缩,释放 {} 个用户数据", reducerIndex, toPersist.size());
}
}
/**
* 将数据写入临时文件
* 真正的生产环境会用Redis或数据库持久化
*/
private void persistToTempFile(int reducerIndex, List<String> userIds,
ConcurrentHashMap<String, OrderSummary> reducer) {
// 伪代码:实际应使用文件IO或Redis
// 这里仅展示思路
String tempKey = "order_summary:reducer:" + reducerIndex;
userIds.forEach(userId -> {
OrderSummary summary = reducer.get(userId);
// redisTemplate.opsForHash().put(tempKey, userId, summary);
});
}
/**
* 最终合并所有Reducer的结果
*/
private Map<String, OrderSummary> mergeReducers(
List<ConcurrentHashMap<String, OrderSummary>> reducers) {
Map<String, OrderSummary> finalResult = new HashMap<>();
for (int i = 0; i < reducers.size(); i++) {
ConcurrentHashMap<String, OrderSummary> reducer = reducers.get(i);
reducer.forEach((userId, partialSummary) -> {
finalResult.merge(userId, partialSummary, OrderSummary::merge);
});
log.info("Reducer-{} 合并完成,当前结果集中用户数: {}", i, finalResult.size());
}
// 计算最终平均值
finalResult.values().forEach(summary -> {
if (summary.getTotalOrders() > 0) {
summary.setAvgAmount(
summary.getTotalAmount()
.divide(BigDecimal.valueOf(summary.getTotalOrders()), 2, RoundingMode.HALF_UP)
);
}
});
return finalResult;
}
}
第四步:组装完整的处理流程
/**
* 订单数据处理入口
* 展示完整的"分-治-合"流程
*/
@Service
public class OrderAggregationService {
@Autowired
private OrderQueryService queryService;
@Autowired
private OrderReducer reducer;
@Autowired
private OrderSummaryRepository summaryRepository;
/**
* 执行订单汇总任务
* 这是整个系统的核心编排方法
*/
public void executeAggregation() {
log.info("========== 开始订单汇总任务 ==========");
long startTime = System.currentTimeMillis();
// Step 1: 获取分页数据流(不一次性加载)
log.info("Step1: 开始分页读取订单数据...");
try (Stream<Order> orderStream = queryService.batchQueryOrders(10000)) {
// Step 2: 并行处理(多核CPU充分利用)
log.info("Step2: 开始Reducer合并处理...");
Map<String, OrderSummary> result = orderStream
.parallel() // 开启并行流,利用多核
.collect(Collectors.groupingByConcurrent(
Order::getUserId,
Collectors.collectingAndThen(
Collectors.toList(),
orders -> buildFinalSummary(orders)
)
));
// Step 3: 写入结果
log.info("Step3: 写入汇总结果到数据库...");
result.values().forEach(summaryRepository::saveOrUpdate);
} catch (Exception e) {
log.error("订单汇总失败", e);
throw new RuntimeException("订单汇总任务执行失败", e);
}
long elapsed = System.currentTimeMillis() - startTime;
log.info("========== 订单汇总完成,耗时: {} 秒 ==========", elapsed / 1000);
}
/**
* 构建用户最终汇总
*/
private OrderSummary buildFinalSummary(List<Order> orders) {
long totalOrders = orders.size();
BigDecimal totalAmount = orders.stream()
.map(Order::getAmount)
.reduce(BigDecimal.ZERO, BigDecimal::add);
long totalItems = orders.stream()
.mapToLong(Order::getItemCount)
.sum();
BigDecimal avgAmount = totalOrders > 0
? totalAmount.divide(BigDecimal.valueOf(totalOrders), 2, RoundingMode.HALF_UP)
: BigDecimal.ZERO;
return OrderSummary.builder()
.userId(orders.get(0).getUserId())
.totalOrders(totalOrders)
.totalAmount(totalAmount)
.totalItems(totalItems)
.avgAmount(avgAmount)
.build();
}
}
四、内存优化效果对比
小陈把这套方案部署后,做了一个详细的监控对比:
┌─────────────────────────────────────────────────────────┐
│ 内存占用对比(峰值) │
├──────────────────┬─────────────┬────────────────────────┤
│ 方案 │ 内存峰值 │ 说明 │
├──────────────────┼─────────────┼────────────────────────┤
│ 旧方案(全量加载) │ 3.2 GB │ 一次性加载全部订单 │
│ │ │ 导致GC频繁,Full GC │
├──────────────────┼─────────────┼────────────────────────┤
│ 新方案(Reducer) │ 0.8 GB │ 分批处理+局部聚合 │
│ │ │ 峰值稳定,无Full GC │
├──────────────────┼─────────────┼────────────────────────┤
│ 节省内存 │ 75% ↓ │ 效果显著 │
└──────────────────┴─────────────┴────────────────────────┘
┌─────────────────────────────────────────────────────────┐
│ 处理耗时对比 │
├──────────────────┬─────────────┬────────────────────────┤
│ 方案 │ 耗时 │ 说明 │
├──────────────────┼─────────────┼────────────────────────┤
│ 旧方案 │ 18分钟 │ 大量GC暂停,实际吞吐低 │
│ │ │ 内存溢出重试2次 │
├──────────────────┼─────────────┼────────────────────────┤
│ 新方案 │ 6分钟 │ 并行处理+稳定内存 │
│ │ │ 吞吐量提升3倍 │
├──────────────────┼─────────────┼────────────────────────┤
│ 性能提升 │ 66% ↓ │ 时间减少 │
└──────────────────┴─────────────┴────────────────────────┘
五、这个方案为什么有效
很多人第一反应是:不就是换个写法吗,怎么效果差这么多?
这里的关键是三个核心设计:
第一个设计:分块输入,而不是全量输入
传统做法: [DB] ────────> [全部加载进内存] ────> [处理]
数据量: 1000万条 内存占用: 3GB 崩溃风险高
Reducer做法: [DB] ─> [块1:1万] ─┐
[块2:1万] ─┤──> [Reducer聚合] ──> [结果]
[块3:1万] ─┘ 内存占用: 0.5GB 稳定运行
...
[块1000:1万]
数据像水流一样,一小股一小股地进来,处理完就流走,不会堆积。
第二个设计:局部聚合,减少中间状态
不是每条数据都创建一个完整的OrderSummary对象。而是在Reducer内部,相同userId的数据先合并成一个”部分结果”,只有相同用户的后续订单才去更新这个部分结果。
// 关键代码:merge操作是增量式合并,不是创建新对象
reducer.merge(userId, newSummary, OrderSummary::merge);
// 注意:merge方法直接修改existing对象,不创建新对象
这就像你记账,不是每花一分钱就新开一本账本,而是在同一行记录上不断追加。
第三个设计:溢出压缩,防止内存爆掉
// 当Reducer缓冲区过大时,主动把小用户数据"推"出去
if (reducer.size() > BUFFER_SIZE) {
compressReducer(reducer, reducerIndex, reducers);
}
这是一个安全阀设计。想象一下,如果某个用户恰好有100万条订单,他的汇总数据会一直留在内存中。通过压缩策略,把大量的小用户数据先持久化,腾出空间给大块数据。
六、给小学生的比喻解释
如果你是小朋友,可以这样理解这个方案:
想象老师要统计全班同学的零花钱总额。
旧方法:让所有同学把零花钱倒在讲台上,老师一个一个数。讲台(内存)放不下,钱撒了一地(内存溢出)。
Reducer方法:把同学分成8组,每组选一个组长。每人把钱交给组长,组长各自统计。最后8个组长把各自的小计交给老师。老师只需要同时拿着8份小计,而不是100份零钱。
这样讲台不会挤爆,老师也不会数不过来。
七、生产环境的注意事项
小陈在实际上线过程中,还遇到了一些坑,分享给你:
坑1:哈希冲突导致数据倾斜
// 问题:如果userId都是数字且分布不均,某些Reducer会特别忙
private int getReducerIndex(String userId) {
return Math.abs(userId.hashCode()) % reducerCount;
}
// 优化:加入盐值,打散哈希分布
private int getReducerIndex(String userId) {
String salted = "order_reduce_salt_" + userId;
return Math.abs(salted.hashCode()) % reducerCount;
}
坑2:并行流的线程安全问题
// 错误:普通HashMap在并行环境下会出问题
Map<String, OrderSummary> result = orders.parallel()
.collect(Collectors.toMap(Order::getUserId, this::buildSummary));
// 正确:使用ConcurrentHashMap或groupingByConcurrent
Map<String, OrderSummary> result = orders.parallel()
.collect(Collectors.groupingByConcurrent(
Order::getUserId,
Collectors.collectingAndThen(
Collectors.toList(),
this::buildFinalSummary
)
));
坑3:超时和断点续传
千万级数据处理不可能一蹴而就,必须支持中断恢复:
/**
* 支持断点续传的Reducer
* 每次处理前,先加载已有的部分结果
*/
@Service
public class ResumableOrderReducer {
private static final String CHECKPOINT_KEY = "order_reduce_checkpoint";
/**
* 加载断点,从上次中断的地方继续
*/
public Map<String, OrderSummary> reduceWithCheckpoint(Stream<Order> orders) {
// 1. 加载上次处理到的进度
Map<String, OrderSummary> checkpoint = loadCheckpoint();
// 2. 跳过已处理的用户
Map<String, Long> processedUsers = checkpoint.keySet().stream()
.collect(Collectors.toMap(
Map.Entry::getKey,
e -> e.getValue().getTotalOrders()
));
// 3. 只处理增量数据
return orders
.filter(order -> !processedUsers.containsKey(order.getUserId()))
.collect(Collectors.groupingByConcurrent(
Order::getUserId,
Collectors.collectingAndThen(
Collectors.toList(),
ordersList -> mergeWithCheckpoint(ordersList, checkpoint)
)
));
}
/**
* 保存断点(每次处理完一批就保存)
*/
public void saveCheckpoint(Map<String, OrderSummary> currentResult) {
// 伪代码:实际应使用Redis或数据库
// redisTemplate.opsForHash().putAll(CHECKPOINT_KEY, currentResult);
log.info("断点已保存,当前已处理 {} 个用户", currentResult.size());
}
}
八、方案的扩展性
这套Reducer架构不只是能处理订单汇总,它可以扩展到很多场景:
场景1: 用户行为日志聚合
── 输入: 每日10亿条点击日志
── Reducer: 按用户ID分片
── 输出: 每个用户的活跃时长、页面访问深度
场景2: 实时风控汇总
── 输入: 每分钟数万笔交易
── Reducer: 按商户ID分片
── 输出: 每笔交易的商户风险评分
场景3: 库存周转分析
── 输入: 各仓库出库流水
── Reducer: 按SKU分片
── 输出: 每个SKU的周转天数
核心思路不变:分而治之,局部聚合,最终合并。
九、最后的一点感悟
小陈在项目复盘时写道:
“很多人一听’千万级数据’就慌了,第一反应是换更贵的服务器、换更重的框架。但实际上,很多时候问题的本质不是数据量大,而是处理方式不对。
Reducer不是一个新技术,它是一个思想——把大问题拆成小问题,逐个解决,最后汇总。这个思想在Excel里能实现,在SQL里能实现,在Java里同样能实现。
技术选型不重要,重要的是理解数据的流动方式。”
这个案例最打动人的地方不在于用了什么高大上的技术,而在于用简单的方法解决了实际问题。没有引入Hadoop,没有搭建Spark集群,就是纯Java,靠合理的架构设计,就把一个内存溢出的问题变成了稳定运行的日常任务。
如果你也在面对类似的数据处理挑战,不妨先从”分块”和”局部聚合”这两个关键词开始想。有时候,最简单的方案才是最有效的。