Skip to content
Charles Shao
Go back

批处理与请求合并:攒批、削频与吞吐优化

views

每发一条日志就做一次网络调用,每插入一行就执行一次 INSERT,每次读取用户名就查一次数据库——这些模式在低流量下毫无感觉,在高流量下却是系统吞吐的天花板。**批处理(Batching)与请求合并(Request Coalescing)**是同一直觉的两种形态:把多次小操作合并成一次大操作,将固定开销摊薄,用可接受的延迟换取更高的吞吐。

本篇从这一直觉出发,讲清客户端攒批、服务端合并、批量读写、异步刷盘的模式与参数,以及延迟-吞吐曲线、背压与常见反模式。

TL;DR

Table of contents

Open Table of contents

1. 核心直觉:固定开销摊薄

任何一次 I/O 操作都有固定开销(Fixed Cost)——无论有效载荷是 1 字节还是 1 MB,这部分成本都要付:

操作固定开销可摊薄部分
网络调用TCP 握手 + RTT(局域网 < 1ms,跨机房 5–50ms)有效载荷的传输时间
数据库写入获取连接 + 事务 begin/commit + 锁实际 SQL 执行时间
磁盘落盘fsync 系统调用(SSD 0.1–1ms,HDD 5–10ms)实际写入数据量
Redis 命令网络 RTT + 命令解析命令本身的计算量

摊薄原理:将 N 次操作合并为 1 次,固定开销只付 1 次;若固定开销占总延迟的比例越高,批处理的收益越显著。

固定开销摊薄:N 次逐条请求 vs 1 次批量请求的代价对比 左侧每条请求独立承担 RTT + 固定开销;右侧 N 条请求只付一次,单条均摊成本降至 1/N。

一句话:批处理不是魔法,它只是把”每次操作都要付的门票钱”变成了”N 个操作共付一张门票”——N 越大,单操作成本越低,吞吐越高,但每个操作等待的时间(延迟)也会增加。

以 JDBC 批量插入为例,固定开销的摊薄效果立竿见影:

// 逐行插入:每行一次网络 RTT + 一次事务 commit
for (LogEvent event : events) {
    stmt.executeUpdate(
        "INSERT INTO event_log(user_id, action, ts) VALUES (?, ?, ?)",
        event.userId, event.action, event.timestamp
    );
}
// 1000 条 × 1ms(RTT + commit) = ~1000ms

// 批量插入:1000 条共用一次网络 RTT 和一次 commit
PreparedStatement ps = conn.prepareStatement(
    "INSERT INTO event_log(user_id, action, ts) VALUES (?, ?, ?)"
);
for (LogEvent event : events) {
    ps.setLong(1, event.userId);
    ps.setString(2, event.action);
    ps.setTimestamp(3, event.timestamp);
    ps.addBatch();
}
ps.executeBatch();    // 一次网络往返,一次 commit
conn.commit();
// 1000 条 × ~0.01ms(均摊) = ~10ms,提升约 100×

连接与 I/O 模型中 Little’s Law 的逻辑呼应:批处理通过增大单次请求的有效载荷,提高了连接利用率,同等连接池大小可以服务更高的总吞吐。

2. 两个触发条件:大小阈值与超时阈值

攒批的触发逻辑由两个条件共同控制,任一满足则立即发送:

SEND if:
  accumulated_size >= batch_size_threshold   // 大小阈值
  OR elapsed_since_first_item >= linger_ms   // 超时阈值(linger)

攒批双触发:大小阈值或超时阈值任一满足即 flush 两个触发条件缺一不可:大小阈值保证高峰充分摊薄,超时阈值保证低峰消息不被饿死。

为什么两者缺一不可

只设大小阈值只设超时阈值
低峰期单个消息永远凑不满一批,长期滞留在缓冲区,消费者看到极高延迟(“饿死”)高峰期大量消息被超时触发,每批都很小,未充分利用批处理收益

参数取舍

参数低延迟场景高吞吐场景
batch_size较小(16–64 KB)较大(256 KB–1 MB)
linger_ms极短(0–5ms)较长(20–100ms)
# Kafka Producer 典型配置示例
# 低延迟场景(广告竞价日志,延迟敏感)
batch.size: 16384      # 16 KB
linger.ms: 1           # 最多等 1ms

# 高吞吐场景(行为埋点批量上报,可接受 100ms 延迟)
batch.size: 1048576    # 1 MB
linger.ms: 100         # 最多等 100ms
compression.type: lz4  # 压缩进一步降低网络开销

3. 客户端攒批

3.1 Kafka Producer 攒批

Kafka Producer 内置攒批机制:消息先写入 RecordAccumulator(按 Partition 分桶的内存缓冲区),由 Sender 线程在满足 batch.sizelinger.ms 时统一发送。

Properties props = new Properties();
props.put("bootstrap.servers", "kafka:9092");
props.put("key.serializer",   "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("batch.size",   "131072");   // 128 KB
props.put("linger.ms",    "20");       // 最多等 20ms
props.put("compression.type", "lz4"); // 压缩减少网络传输
props.put("acks",         "all");      // 全副本确认,不丢消息

KafkaProducer<String, String> producer = new KafkaProducer<>(props);
// send() 调用是异步的,消息先入 RecordAccumulator,由 Sender 线程攒批发出
producer.send(new ProducerRecord<>("events", key, value));

3.2 Redis Pipeline

Redis 单命令执行包含一次网络 RTT;Pipeline 将多条命令打包为一次网络往返,服务端顺序执行后一次性返回所有结果:

// 不用 Pipeline:N 条命令 × N 次 RTT
for (String key : keys) {
    String value = jedis.get(key);  // 每次 get 都是一次网络往返
}

// 使用 Pipeline:N 条命令只有 1 次 RTT
Pipeline pipeline = jedis.pipelined();
Map<String, Response<String>> responses = new LinkedHashMap<>();
for (String key : keys) {
    responses.put(key, pipeline.get(key));  // 命令先入队列
}
pipeline.sync();  // 一次性发送并等待所有响应
responses.forEach((key, resp) -> process(key, resp.get()));

注意:Pipeline 不是原子操作(与事务 MULTI/EXEC 不同),中间某条命令失败不会回滚其他命令;也不保证所有命令在同一个 Redis Cluster 节点上执行(跨槽命令需手动分组)。

3.3 日志与埋点上报

客户端 SDK 的日志上报通常采用内存缓冲 + 定时批量发送的模式:

public class EventReporter {
    private final BlockingQueue<Event> buffer = new LinkedBlockingQueue<>(10000);
    private final ScheduledExecutorService scheduler;

    public EventReporter() {
        scheduler = Executors.newSingleThreadScheduledExecutor();
        // 每 500ms 或缓冲满 200 条时触发一次批量上报
        scheduler.scheduleAtFixedRate(this::flush, 500, 500, TimeUnit.MILLISECONDS);
    }

    public void report(Event event) {
        if (!buffer.offer(event)) {
            // 缓冲满:丢弃(埋点允许少量丢失)或触发同步上报
            dropOrFlushSync(event);
        }
    }

    private void flush() {
        List<Event> batch = new ArrayList<>(200);
        buffer.drainTo(batch, 200);          // 取出最多 200 条
        if (!batch.isEmpty()) {
            httpClient.postBatch("/events", batch);
        }
    }
}

关键参数:缓冲队列大小(防止内存无限增长)、批大小上限(防止单次请求过大)、flush 间隔(控制最大延迟)。

3.4 客户端攒批横向对比

场景工具 / 机制大小参数超时参数备注
Kafka 生产者RecordAccumulatorbatch.sizelinger.ms内置;压缩配合使用收益更大
Redis 多命令Pipeline手动控制批数手动控制触发时机非原子;Cluster 需按槽分组
JDBC 写入addBatch / executeBatch手动控制批数定时 flush需在事务内;批大小建议 500–2000
HTTP 上报内存缓冲 + 定时 flush条数 / 字节上限flush 间隔适合埋点、日志、监控指标

4. 服务端请求合并(Singleflight)

客户端攒批解决的是”我主动把多次调用合并”;Singleflight 解决的是”同一时间多个调用者并发请求同一资源时,只让一个请求真正去取,其余等待并共享结果”。

这与缓存机制缓存击穿的解法本质相同:热点 key 过期后,大量并发请求同时 miss,若每个都去数据库取,会造成瞬间冲击——Singleflight 把”同一 key 的 N 个并发回源”折叠为 1 次。

// 使用 Guava 的 LoadingCache 或自行实现 Singleflight
// 以下用 ConcurrentHashMap + CompletableFuture 示意
private final ConcurrentHashMap<String, CompletableFuture<Product>> inflightRequests
    = new ConcurrentHashMap<>();

public Product getProduct(String productId) throws Exception {
    CompletableFuture<Product> future = inflightRequests.computeIfAbsent(
        productId,
        key -> CompletableFuture
            .supplyAsync(() -> loadFromDB(key))  // 只有第一个请求真正回源
            .whenComplete((v, e) -> inflightRequests.remove(key))  // 完成后移除
    );
    return future.get();  // 后续并发请求复用同一个 future
}

Singleflight 请求合并:同一 key 的并发请求折叠为一次回源,结果广播给全部等待者 N 个并发请求中只有第一个穿透到数据库;其余等待同一个 Future,回源完成后同时获得结果。

适用场景

不适用场景:写操作或有副作用的操作——Singleflight 假设相同参数的调用结果完全等价,写操作不满足此前提。

5. 批量读写

5.1 批量查询:消灭 N+1

N+1 问题是最常见的数据库性能反模式:先查出 N 个对象的 ID,再逐个查询其关联数据,产生 N+1 次数据库调用。

-- 反模式:N+1 查询
SELECT id FROM orders WHERE user_id = 123 LIMIT 20;
-- 然后对每个 order_id 执行:
SELECT * FROM order_items WHERE order_id = ?;  -- 执行 20 次

-- 优化方式一:JOIN(适合关联数据量小时)
SELECT o.*, oi.*
FROM orders o
JOIN order_items oi ON oi.order_id = o.id
WHERE o.user_id = 123 LIMIT 20;

-- 优化方式二:批量 IN 查询(适合关联数据量大或需要分表路由时)
SELECT * FROM order_items
WHERE order_id IN (101, 102, 103, ..., 120);  -- 一次查询 20 个

DataLoader 模式:在 GraphQL 场景或需要跨多个数据源批量加载的场景,DataLoader 把同一请求周期内(一个 event loop tick)的所有单个查询收集起来,攒成一次批量查询:

// 每个 resolveUser(userId) 调用都会被 DataLoader 收集
// 在 tick 结束时合并为 userLoader.loadMany([1, 2, 3, ...])
DataLoader<Long, User> userLoader = DataLoader.newDataLoader(
    (List<Long> userIds) -> userService.findAllByIds(userIds)  // 一次批量查询
);

5.2 批量写入

JDBC batch insertbulk upsert 是减少写入 RTT 的标准方式:

-- 批量 INSERT(MySQL 支持 VALUES 多行写入)
INSERT INTO events (user_id, action, ts)
VALUES
    (1001, 'click', '2026-07-13 10:00:00'),
    (1002, 'view',  '2026-07-13 10:00:01'),
    (1003, 'buy',   '2026-07-13 10:00:02');
-- 相比 3 次单行 INSERT:节省 2 次网络 RTT 和 2 次事务 commit

-- 批量 UPSERT(MySQL ON DUPLICATE KEY UPDATE)
INSERT INTO user_stats (user_id, click_count)
VALUES (1001, 5), (1002, 3)
ON DUPLICATE KEY UPDATE click_count = click_count + VALUES(click_count);

批大小建议:单批 500–2000 行是较合理的起点;超过 5000 行时单批事务持锁时间变长,影响并发写入;具体上限通过压测确定(参见性能测试)。

5.3 mget / hmget

Redis 的 mgethmget 命令原生支持批量 key 读取,与 Pipeline 相比,它们是单条命令(一次 RTT,服务端原子批量读),无需手动管理 Pipeline 的分组问题:

// 单次 mget 取 N 个 key,只有 1 次 RTT
List<String> values = jedis.mget("user:1001", "user:1002", "user:1003");

// hmget 批量取同一个 Hash 的多个字段
List<String> fields = jedis.hmget("user:1001", "name", "email", "score");

Redis Cluster 注意mget 的多个 key 可能落在不同槽(Slot),Cluster 客户端(如 Lettuce、Jedis Cluster)会自动按槽分组发到对应节点,但并行 RTT 变为多次,需权衡。通过 Hash Tag({user}:1001{user}:1002)可强制多个 key 落同一槽。

6. 异步刷盘与组提交(Group Commit)

6.1 WAL 与 fsync 成本

数据库写入的持久化依赖 WAL(Write-Ahead Log):数据先写入 WAL,fsync 刷盘后事务才能提交。fsync 是最昂贵的系统调用之一——它强制将操作系统 Page Cache 中的数据写入磁盘:

6.2 组提交(Group Commit)

组提交将多个并发事务的 WAL 刷盘合并为一次 fsync

事务 T1 写入 WAL → 等待 fsync
事务 T2 写入 WAL → 等待 fsync    ┐
事务 T3 写入 WAL → 等待 fsync    ├─→ 一次 fsync 刷盘 → T1、T2、T3 同时提交
事务 T4 写入 WAL → 等待 fsync    ┘

效果:并发写入越高,组提交的收益越显著;单机 TPS 可从”1 次 fsync = 1 个事务”跃升到”1 次 fsync = N 个事务”。MySQL InnoDB 默认启用 Binlog Group Commit(MySQL 5.6+),Redo Log 的 Group Commit 由 innodb_flush_log_at_trx_commit 控制:

-- innodb_flush_log_at_trx_commit = 1(默认):每次提交都 fsync,最安全,性能最低
-- innodb_flush_log_at_trx_commit = 2:每次提交写 OS 缓存,每秒 fsync 一次,宕机丢 1s 数据
-- innodb_flush_log_at_trx_commit = 0:每秒写缓存并 fsync,吞吐最高,宕机丢 1s 数据

对于非金融类写入(如埋点、日志),设置为 2 可显著提升写吞吐,接受极小的数据丢失窗口。

6.3 异步刷盘的边界

异步刷盘(fsync 延迟或不做 fsync)是”吞吐 vs 持久性”的显式取舍:

配置吞吐宕机丢失风险适用场景
同步 fsync(每次提交)金融核心交易、账务数据
每秒 fsync中高最多 1s 数据日志、埋点、统计计数
fsync(纯内存写)最高宕机丢失全部未刷数据缓存、可重算的中间结果

一句话:组提交和异步刷盘本质上是在”我愿意等多久”和”我愿意在宕机时丢多少”之间的连续取舍——选择之前先想清楚业务对持久性的真实要求。

7. 延迟-吞吐权衡曲线

7.1 linger 与延迟-吞吐的关系

攒批引入的延迟是确定性的:一条消息/请求从进入缓冲区到被发送,最长等待时间等于 linger_ms。因此:

P99 延迟下界 ≈ 业务处理时间 + linger_ms + 网络 RTT
吞吐上界   ≈ batch_size / (业务处理时间 + linger_ms + 网络 RTT)

如何设 linger

场景推荐 linger理由
广告竞价日志(DSP 超时 100ms)1–5ms日志上报不在关键路径,但延迟预算极紧
用户行为埋点(分析用)50–200ms对延迟不敏感,吞吐优先,节省带宽和服务端处理压力
数据库批量写入(异步队列消费)10–50ms,或积累 N 条触发吞吐优先,可接受秒级端到端延迟
实时支付流水0ms(禁用批处理)强一致、低延迟,批处理带来的等待不可接受

7.2 背压与内存上限

攒批缓冲区是内存中的队列,若下游发送速度跟不上上游写入速度,缓冲区会持续增长,最终 OOM:

生产速率 > 消费速率 → 缓冲区增长 → OOM

背压策略

Kafka Producer 的 buffer.memory(默认 32MB)和 max.block.ms(默认 60s)是这一机制的生产实现:缓冲满时 send() 调用阻塞,超过 max.block.ms 后抛出 TimeoutException

8. 反模式与小结

8.1 常见反模式

反模式一:批太大导致长尾延迟 / OOM

batch.size 设为 64MB,单批写入包含数万行——数据库单事务持锁时间极长,P99 延迟飙升;网络传输中出现一次重传就可能超时;一批失败时需要整批重试,放大了重试开销。

批大小应与下游能力和超时预算对齐:先用较小的批大小上线,通过压测逐步找到吞吐-延迟的最优点,而不是拍脑袋设一个”大一点没关系”的值。

反模式二:只设大小阈值,不设超时阈值

linger.ms = 0batch.size = 1MB——低峰期流量稀少,消息积累缓慢,数十分钟也凑不满 1MB,下游消费者长时间收不到数据。超时阈值是低峰期的”保底发送”机制,必须与大小阈值同时设置。

反模式三:批处理里一条失败拖垮整批

批量 INSERT 中第 500 条违反唯一约束——整批事务回滚,前 499 条都需要重试,且重试时仍然可能在第 500 条失败,进入死循环。

正确处理方式:

try {
    ps.executeBatch();
    conn.commit();
} catch (BatchUpdateException e) {
    int[] updateCounts = e.getUpdateCounts();
    List<Integer> failedIndices = new ArrayList<>();
    for (int i = 0; i < updateCounts.length; i++) {
        if (updateCounts[i] == Statement.EXECUTE_FAILED) {
            failedIndices.add(i);
        }
    }
    // 对失败的条目单独处理(记录死信、跳过、告警),其余条目正常提交
    retryOrDeadLetter(failedIndices, batch);
}

或改用 INSERT IGNORE / ON DUPLICATE KEY UPDATE 让数据库层面静默处理冲突,避免整批回滚。注意:批量写”部分成功”是分布式写入的常态,必须配合幂等(唯一约束 / event_id 去重)才能安全重试,否则重复插入在计费场景就是重复扣费。

反模式四:批处理遮蔽了真实的业务延迟

将本应实时响应的用户操作(如下单确认)异步丢入批处理队列——用户等待 200ms 后才收到确认,体验极差。批处理适用于异步链路(日志、埋点、数据同步),不适用于同步响应的关键路径。

8.2 小结

批处理与请求合并是”用延迟换吞吐”的通用手段,适用于链路中每一处存在固定开销的地方:

结合异步与削峰(控制流入速率)、连接与 I/O 模型(高效利用出向连接)、数据库性能与扩展(批量写入的落地目标)、幂等与一致性(保证部分失败可安全重试)与性能测试(验证批大小和 linger 的实际效果),批处理是高性能系统中最轻量、收益最确定的优化手段之一——往往只需改几行配置,就能在不换机器的前提下把吞吐翻倍。它既能扛住数十万 QPS 的写入,也会因一个 linger 的误设而 OOM——参数即代码,务必校验、监控、审查。


views
Share this post on:

Previous Post
性能测试:指标、负载模型、结果解读与容量规划
Next Post
连接与 I/O 模型:线程、池化、多路复用与背压