RecordAccumulator是Kafka生产者的消息缓冲区,负责攒批、复用内存(BufferPool管理16KB ByteBuffer)、按TopicPartition分队列存储,由Sender线程批量发送,避免频繁GC,提升吞吐。

缓冲区总容量受 buffer.memory 控制
默认 32MB(33554432 字节),由参数 buffer.memory 配置。所有分区共用这块总内存。当 append 新消息时,若剩余可用内存不足,且未设置 max.block.ms=0,线程会阻塞等待;设为 0 则直接抛出 `TimeoutException`。
按 batch.size 复用固定大小的 ByteBuffer
默认 16KB(16384 字节),由参数 batch.size 决定。`RecordAccumulator` 内部维护一个 BufferPool,只管理恰好等于 batch.size 的 ByteBuffer 实例:
- 每次需要新建 ProducerBatch 时,优先从 BufferPool 取一个空闲的 16KB buffer
- ProducerBatch 被发送完成或被丢弃后,其 backing buffer 会被归还给 BufferPool,而非直接 GC
- 如果申请的 buffer 大小 ≠ batch.size(比如压缩后超限或手动调大单条消息),则绕过 BufferPool,直接用
ByteBuffer.allocate()分配,这类 buffer 不回收
每个分区独占队列,按需动态创建批次
`RecordAccumulator` 使用 ConcurrentMap<TopicPartition, Deque<ProducerBatch>> 组织数据:
在 Java 中初始化和管理阿里云 SDK客户端。包括单例模式、线程安全、endpoint 与 region 配置、VPC 终端节点、同步与异步等。
- 每个 Topic-Partition 对应一个双端队列(Deque),保证该分区内的消息严格有序
- 新消息先尝试追加到队尾的“当前活跃批次”(last batch)中;成功则复用该 batch 的 buffer
- 若 last batch 已满、已关闭或不存在,则触发 allocate → 从 BufferPool 拿 buffer → 创建新 ProducerBatch → 入队
- 队列本身不占 buffer.memory,只存对象引用,内存压力集中在 ProducerBatch 的 MemoryRecords 所持有的 ByteBuffer 上
内存释放依赖 Sender 线程驱动
BufferPool 中的 buffer 只有在 ProducerBatch 被确认发送成功、失败重试放弃、或生产者关闭时才会真正归还:
立即学习“Java免费学习笔记(深入)”;
- Sender 线程将批次发往 Broker 后,收到响应(ack)或判定超时/异常,会调用
completeBatch()或abortBatch() - 这些方法内部触发
done()→ 清理 metadata → 将 buffer 返还 BufferPool - 若批次因 linger.ms 触发发送,或因 batch.size 满而关闭,只要还没交到 Sender 手里,buffer 就一直被持有


















