尧图精选

数据流式编程中的堆数据堆积与原地重用优化实践

🕒 发布时间:2026/9/15 2:17:26 📁 来源:尧图网络
先说个我自己的经历。前阵子排查一条跑在数据流式编程框架上的处理链路每秒要吞二十多万条设备心跳消息业务逻辑不算复杂无非是解析、补字段、聚合、落库。诡异的是不管怎么压测吞吐始终上不去CPU 并不高内存却像漏了一样GC 日志里 Minor GC 每秒钟来好几回。把堆 dump 拉下来一看堆里躺着几百万个一模一样的小对象全是处理过程中新建的中间数据结构。那一刻我意识到一个很朴素但很深刻的问题很多流式任务真正的性能瓶颈根本不是算法复杂度而是堆数据在生产线上不断“制造垃圾”。而解决这个问题的核心思路说穿了就是标题里的“原地重用”——不新建只覆写。今天这篇文章就围绕这三个关键词展开数据流式编程出现的堆数据堆积问题、为什么“原地重用”是最直接的解法、以及我实际落地时的实现方式、测试数据和踩坑记录。想优化流式处理链路的同学这篇应该能帮你省下不少弯路。1. 数据流式编程与堆数据为什么这对组合这么容易出问题1.1 流式模型的本质数据不动算子动理解流式编程之前先放下框架层面的概念回到最朴素的模型。数据流式编程stream-based programming本质上是一条生产流水线数据从源头进来经过一个又一个处理节点最后到达终点。每个节点接收上游数据、做处理、把结果交给下游。这个模型里数据本身不像传统批处理那样被整体收集起来再慢慢算而是“边到边走”像水流一样持续流动。流式模型最大的优势是延迟低、内存占用有上界。听起来很美但工程实现里有个隐蔽的问题为了让数据在节点之间传递每个算子往往都要产出新的数据对象这些对象就是堆上的动态分配。比如一个 filter 节点上游传进来一个事件它判断完条件后把事件继续往下传时如果框架的数据模型要求“不变性”那它很可能复制一份再传。这一步一次分配堆上就多一个对象。我见过不少刚接触流式处理的团队把精力全放在调窗口大小、调并行度上却没人注意处理链路里每秒钟究竟在堆上分配了多少对象。说实话这个数字往往比想象中大一到两个数量级。1.2 不可变风格的代价每次转换都是一次堆分配这里要展开说说“不变性”这个东西。很多流式框架、函数式风格的代码里数据对象被设计成不可变的方便并发安全、方便推断逻辑。但不可变对象一旦要“变化”就只能新建。举个例子一个简单的 map 操作想把事件里的温度从摄氏度转成华氏度不可变设计下你就得 new 一个新的事件对象把原对象字段拷过去改一下温度字段。这一趟下来堆上多了一个对象多了一次内存拷贝。如果这是偶发操作完全无所谓。可流式链路是每时每刻都在跑的一秒二十万条消息就意味着每秒二十万个短命对象诞生。它们很快变成垃圾等着被回收。短命对象恰恰是新生代 GC 的主要压力来源分配快回收也快但哪怕回收再快次数一多STWStop-The-World停顿照样把吞吐吃掉一大截。我在那次排查里看到的堆 dump 很有代表性几百个 Live 对象里超过九成是各种中间态——Event 拷贝、ArrayList 的 elementData 数组、各种迭代器和包装器。业务大头反而只占一小块。这就是典型的“处理逻辑没多少对象制造机开得飞起”。1.3 不只是 GC缓存和内存碎片也在拖后腿堆数据堆积的影响不只是 GC。还有两个层面经常被忽略一个是 CPU 缓存局部性一个是内存碎片。先说缓存。当你每次新建对象时对象在堆里的物理地址不一定连续。第一次创建的对象可能在内存地址 A第二次的却跑到地址 D中间隔了好远。CPU 缓存是按缓存行加载的一般是 64 字节如果数据分布跳来跳去缓存命中率就会很差。相反如果你复用同一个对象它的地址始终不变相关字段一直留在缓存行里每次访问都是“热”的处理速度自然更快。再说内存碎片。频繁地分配、释放大小不一的对象很容易把堆内存切得零零碎碎。JVM 和不少运行时都有分配缓冲区机制来缓解但分配压力过大时碎片化依然会让分配变慢甚至提前触发 Full GC 做堆整理。很多线上案例里“明明内存还很充足却频繁 Full GC”往往就是碎片化导致的。所以数据流式编程里的堆数据问题不是某一个点出了问题而是分配压力、GC 停顿、缓存失配、碎片化四条线同时在拖累系统。要治本就不能指望调几个 JVM 参数得回到代码层面减少不必要的堆分配——这正是“原地重用”要解决的问题。2. 原地重用的核心思路复用内存而不是重建对象2.1 什么是堆数据的原地重用“原地重用”这个名字听起来有点唬人说白了就是一句话一块已经分配好的堆内存用完之后不要急着还给 GC而是清一清、改一改让下一个数据接着用它。传统写法里数据对象用完了就丢掉让 GC 去回收需要新对象时再分配。原地重用的写法里你用的是一个可变的“容器”对象每次处理新数据就把它里面旧的值清掉、写入新值再交给下游。整个生命周期里这个对象在堆上的地址不变变的只是它内部的数据。打一个生活化的比方传统写法是去便利店买一瓶新水喝完把瓶子扔了下次渴了再买一瓶。原地重用是你留着一个瓶子喝完接自来水反反复复用。省下的是瓶子的“制造”和“回收”成本代价是你得负责把瓶子洗干净——对应到代码里就是处理好对象的重置逻辑别让上一批数据的残留污染下一批。2.2 三种落地形态对象池、覆写式更新、复用容器我自己的实践中原地重用大致有三种落地形态覆盖面比较广可以根据场景灵活选。第一种是对象池。这个最经典也最通用。预先创建一批对象放进池子处理流里的每一条数据时从池子里借一个对象出来用用完了再还回去。池子里的对象反复被借用、归还从分配器的角度看堆上始终只存在固定数量的一批实例。适合那些创建成本高、复用价值大的对象比如带大数组、大缓冲的数据包或事件对象。第二种是覆写式更新。这种形态不依赖“池”这种显式的容器而是直接把对象设计成可变的处理逻辑每次拿到新数据就调用对象的 setter 或 reset 方法把字段改掉。常见于聚合状态、计数器、累加器这类场景。比如你要算一个窗口内所有事件的总和你不需要每来一条事件就新建一个聚合结果对象维护一个可变的长整型累加变量就够了。第三种是复用容器与缓冲。这里针对的是流式处理中最常见的中间容器List、Map、数组。每次处理一批数据往往要新建一个容器来装中间结果这种分配的代价非常明显。合理的做法是提前分配好容器每次处理时清空再填充。容器底层的大数组直接保留清空只改 size下次添加元素时就不会触发扩容也不用重新分配底层存储。三种形态并不互斥实际项目里经常混着用。比如处理链路上用对象池复用事件对象聚合节点用覆写式更新维护状态窗口计算里用复用容器来存批处理结果。一套组合拳下来堆分配次数能锐减一个数量级以上。2.3 适用性判断什么场景值得做什么场景别碰原地重用不是银弹盲目使用反而会把代码搞得一团糟。我给自己定了几条判断标准供你参考。值得做的场景通常具备这么几个特征第一对象创建频率极高每秒成千上万次甚至更高第二对象生命周期极短从创建到变成垃圾只隔了几个操作第三创建或销毁成本不低比如涉及大数组、复杂初始化、IO 资源第四对象只在线程内部传递不跨线程、不入队列。如果全中优先考虑原地重用。反过来说有些场景真的别碰。如果对象要跨线程传递或者要放进队列等异步消费者处理那对象的生命周期就不受当前线程控制了你这一端刚把对象放回池子另一端可能还在用它数据就错乱了。还有如果是串并行度不高的低频逻辑比如一天只跑几次的批清理任务原地重用带来的性能提升微乎其微却会牺牲代码可读性完全得不偿失。一句话总结我的经验性能收益取决于分配量分配量取决于频率和体积。频率高、体积大的对象优先考虑原地重用频率低、体积小的对象直接用传统写法更省心。3. 实操实录用一个行情处理链路讲透原地重用的每个细节3.1 场景设定事件解析与聚合链路为了让“原地重用”不停留在概念层面我拿之前那个优化过的设备消息处理场景换个壳简化成一套行情事件处理链路展示具体怎么改。假设我们要处理一条行情事件流上游源源不断地送来行情 tick 数据每个 tick 包含时间戳、证券代码、最新价、成交量四个字段。我们的处理逻辑分三步第一步把原始的二进制数据解析成事件对象第二步按证券代码分组累计每个代码当日的成交额第三步把聚合结果周期性地输出。传统的不可变风格代码大概长这样record Tick(long ts, String code, double price, long volume) {} public final class OldPipeline { public void onBinaryData(byte[] data) { Tick tick parseTick(data); // 每次反序列化都新建一个 Tick MapString, MutableAggr aggr aggrMap.get(tick.code()); if (aggr null) { aggr new MutableAggr(); aggrMap.put(tick.code(), aggr); } aggr.add(tick.price(), tick.volume()); // 同样存在重复创建 } }这段代码在每秒二十万条的流量下每秒至少要新建二十万个 Tick 对象。而实际上 Tick 对象只是为了传那几个字段它的生命周期从头到尾不过几微秒。这就是最典型的“短命堆对象”也是我们优化的突破口。3.2 第一步把反序列化改成“解析进复用对象”首先改造的是解析环节。我们不 new Tick而是维护一个可变的 MutableTick 对象每次解析直接把结果写进去final class MutableTick { long ts; String code; double price; long volume; void reset() { ts 0; code null; price 0.0; volume 0; } } public final class PooledPipeline { private final ThreadLocalMutableTick tickRef ThreadLocal.withInitial(MutableTick::new); public void onBinaryData(byte[] data) { MutableTick tick tickRef.get(); parseInto(data, tick); // 解析结果直接写入 tick 内部字段 MutableAggr aggr aggrMap.get(tick.code); if (aggr null) { aggr new MutableAggr(); aggrMap.put(tick.code, aggr); } aggr.add(tick.price, tick.volume); } }这里有个关键点因为整个处理过程都发生在同一个线程内我们可以用 ThreadLocal 直接持有这个复用对象。每个线程最多只有一个 MutableTick 实例线程之间天然隔离不用担心并发问题。而 parseInto 方法内部的写法也要相应调整——不是返回新对象而是把数据填进传入的入参里。有人可能会问String code 不也得新建吗确实反序列化字节数组里提取字符串时通常没法完全避免创建 String。但我们可以做一层优化如果代码是定长的比如 6 位数字证券代码可以直接把它当作 char[] 存在 MutableTick 里而不是用 String比较时用 Arrays.equals 比对字符数组map 的 key 则用自定义的轻量包装对象。这一步能再省掉一大块 String 创建压力。3.3 第二步聚合器设计成覆写式更新刚才代码里出现的 MutableAggr 其实就是第二种形态“覆写式更新”的体现。传统写法里每来一个事件都创建一个新的聚合结果然后把旧值加进去或者用 ConcurrentHashMap 的 merge 方法底层也会涉及多个内部节点的创建。覆写式写法就是直接维护一个可变累加器把聚合字段作为对象属性直接在原对象上累加final class MutableAggr { double totalTurnover; long totalVolume; void add(double price, long volume) { totalTurnover price * volume; totalVolume volume; } void reset() { totalTurnover 0.0; totalVolume 0; } }注意这里的关键不只是“可变”而是“这个可变对象从头到尾只有一份”。如果每来一条事件就 new 一个 MutableAggr 再累加那和之前没区别。真正的收益来自复用同一个累加器实例让数据在字段上原地更新。对于按 key 分组的聚合对象池不急着做。更优雅的方案是维护两个 map一个保存当前周期正在累加的“活跃聚合器”一个保存上一周期已经完成、可以被重置复用的“空闲聚合器”。当前周期结束时把活跃聚合器整体迁移到空闲集合里下个周期来了再从空闲集合里取出来复用。这种思路在 Flink 的窗口状态管理、Kafka Streams 的状态存储里都能看到类似的影子。3.4 第三步窗口数据用环形缓冲复用最后一步是窗口计算。假设我们要算最近 1 分钟内每分钟的成交额传统做法是维护一个 List 新数据来了 add窗口滑动了就把头部元素移除。List 会频繁触发扩容和数据搬移移除头部元素又要批量移动后面的元素。这里的堆分配主要来自底层数组的扩容和反复搬移。我改成了环形缓冲ring buffer来实现这个滑动窗口。环形缓冲用固定大小的数组存储数据用 head 和 tail 两个指针标记窗口边界新数据写到 tail 位置tail 加一窗口溢出时 head 加一相当于把最旧的数据覆盖掉。整个过程中数组本身不变不存在扩容和搬移final class RingBuffer { private final double[] prices; private final long[] volumes; private final long[] timestamps; private int head; private int tail; RingBuffer(int capacity) { prices new double[capacity]; volumes new long[capacity]; timestamps new long[capacity]; } void add(long ts, double price, long volume) { timestamps[tail] ts; prices[tail] price; volumes[tail] volume; tail (tail 1) % timestamps.length; if (tail head) { head (head 1) % timestamps.length; } } }环形缓冲真正做到了“原地”的极致不仅对象不新建连底层数据都是直接覆盖旧值。窗口统计时只需遍历 head 到 tail 之间的有效数据即可。整个窗口内部不产生任何堆分配唯一的代价是初始化时一次性分配三个数组。如果窗口本身很宽比如几千条数据这部分固定内存完全可接受。顺便说一句如果你用 Java 的 ArrayDeque 或者 LinkedList 做滑动窗口底层其实也在频繁分配节点或扩容数组。换成环形缓冲后GC 压力立刻小一大截。这也是流式处理框架底层热衷用环形缓冲的原因之一。3.5 实现细节里的五个隐性坑使用原地重用的实现看起来简单落地时细节不少。我单独整理了几条实战注意点。第一重置方法一定要彻底。字段漏了任何一个上一批数据的残留值就会混进下一批计算。我吃过一次亏某次优化一个带默认值的聚合器忘了重置一个布尔字段结果所有后续窗口的“是否有效”标记全是上一批数据留下的 true排查了很久才发现。建议写一个统一 reset() 方法每次“还回”对象时调用并且用 assert 或者 debug 日志校验关键字段已恢复初始值。第二对象池的归还时机必须非常明确。基于流式处理的代码天生是“链式调用”的对象在节点之间传递时要定义清楚“谁负责归还”。我的做法是约定当前节点把对象传给下游之前不再持有引用下游用完必须归还如果没有下游终点节点则由终点节点归还。这个约定一定要写进团队规范不然后面接手的人会在池里发现各种“霰弹对象”。第三优先用 ThreadLocal 而不是全局池。多数流式处理节点是单线程顺序执行的ThreadLocal 就能满足复用需求而且比全局对象池更简单、更安全也没有锁竞争的问题。只有明确存在多线程交替处理时才考虑全局池并且要配上原子操作或锁。第四复用对象的引用不能逃逸出流处理线程。也就是前面说的“生命周期不受当前线程控制”的对象不能复用。如果你把对象塞进异步队列、提交到线程池、或者存到全局缓存里就必须立刻改为传统的新建模式。原地重用和异步边界天然冲突不要硬融合。第五测性能时别只看 GC。原地重用带来的提升一部分来自 GC 次数减少另一部分来自缓存命中率的提升。跑 benchmark 时除了记录 GC 次数还要看 P99 延迟、CPU 缓存失配率可以用 perf stat -e cache-misses 之类的工具。我测试时明显感觉到GC 停顿降下去之后P99 延迟的改善比平均延迟更明显这对流式系统来说才是关键指标。4. 实测数据原地重用到底能省多少4.1 测试环境与对照组设计改造完成之后我搭了一个简单的基准测试来量化收益。测试环境是我们内部的一台 8 核 16G 的压测机JDK 17堆设了 4G。生产的输入是模拟的行情二进制消息大小约 80 字节一条总共灌入 500 万条数据。对照组是前面提到的 OldPipeline传统不可变风格每次解析都新建 Tick 对象聚合用的 Map 使用 computeIfAbsent 动态创建聚合器。实验组是 PooledPipeline使用 ThreadLocal 持 MutableTick覆写式聚合器环形缓冲做滑动窗口。每组跑 5 次取中位数。统计口径包括总耗时、每秒吞吐、堆分配总量、Minor GC 次数、GC 总时间。这里要说明我们的场景偏计算密集生产环境的网络 IO 可能会稀释这部分收益但压力测试能看出纯计算路径的真实差距。4.2 测试结果与数据解读测试结果整理成了表格我直接贴出来指标传统写法OldPipeline原地重用PooledPipeline提升幅度总耗时ms42312187约 48.3%吞吐万条/秒118.2228.6约 93.4%堆分配总量GB8.60.6约 93.0%Minor GC 次数385约 86.8%GC 总时间ms78294约 88.0%P99 延迟us386142约 63.2%堆分配总量那一列是最惊人的从 8.6G 掉到 0.6G。要知道我们处理的原始数据一共才 500 万条 × 80 字节 400MB 左右。传统写法堆分配了 8.6G说明每 1 字节输入数据平均会产生约 20 字节的垃圾对象重复分配的问题一目了然。原地重用的版本里剩下的 0.6G 基本都是解析时不可避免的字节数组处理和少量初始化开销。吞吐提升将近一倍是因为 GC 停顿大幅减少CPU 时间更多用在了业务计算上。而 P99 延迟的改善也符合预期——少了很多停顿和缓存未命中极端情况下的延迟自然就降下来了。4.3 什么时候差距不明显甚至更慢虽然上面的数据很好看但我必须如实说有几类场景下原地重用的收益会很小甚至出现负优化。第一类是 IO 密集型链路。如果处理逻辑里充斥着网络读取、磁盘写入、下游 RPC 调用堆分配的占比会被 IO 等待大幅稀释。比如一条消息处理耗时 100ms其中 99ms 在等下游响应那你把堆分配降一个数量级整体收益也不到 1%。这种场景优先优化 IO 路径而不是底层对象创建。第二类是分配量本来就很少的场景。有些流式处理的算子是无状态且极简的入参出参都是基本类型压根没有中间对象也就没有“原地重用”的空间。硬套对象池反而增加代码复杂度收益趋近于零。第三类是对象生命周期无法受控的场景。前面反复提过如果处理链路里存在异步边界或者对象需要被下游引用超过当前处理周期原地重用就没法安全地实现。这种情况下强行复用可能会导致数据串号、并发冲突需要加锁或拷贝反而比传统写法更慢。所以做这类优化前一定要先做剖析profiling看清楚了分配热点再动手。不要凭着“直觉”把所有对象都池化一遍那只是把复杂度从一个地方挪到另一个地方而已。5. 常见问题与排查技巧实录5.1 问题一复用对象被意外共享数据“串号”这是我见过最多的问题。现象是处理结果时而正确、时而错误字段偶尔出现上一次事件的残留值并发量一高就频繁出现。原因通常有两种一种是不小心把复用对象放进了全局容器或异步队列多个线程或前后多个处理周期共享了同一个实例另一种是对象归还池子后调用方又持有了它的“旧引用”在后面某处又读了一次旧引用而这时对象已经被其他人改写了。排查时可以给池里的对象加一个 lease/占用标记借用时标记为占用归还时标记为空闲归还后如果再有代码尝试访问这个对象直接抛异常。这个“强制协议”能很快暴露问题点。另一种方法是开启对象 header 里的身份哈希identity hash code日志跟踪对象地址的分配和释放链路。5.2 问题二忘记重置导致的脏数据这个坑前面提过我再单独展开一下。它比较容易出现在两个位置一个是对象放入池子之前另一个是对象每次被借用出来之后。前者不做好脏数据残留到下一次使用后者不做好第一次使用没问题第二次就会用上残留值。我的做法是双保险归还时调用一次 reset()借用时再调用一次 reset()。虽然看起来多了一次调用但 reset 本身只赋值基本类型开销可以忽略。而且这种对称设计符合“资源获取即初始化”的思路能兜住人脑的疏忽。另外reset 方法建议不要只设默认值还要给每个字段一个明确语义。比如 0 表示无数据null 表示无引用。这样出问题时日志里能直观看出哪里没重置干净。5.3 问题三对象池内存泄漏池化之后内存泄漏的“症状”变了堆里对象的数量稳定不涨但内存占用持续走高。很多人会懵其实这是把对象池放得太大造成的。比如池里保留了 1 万个事件对象每个对象内部有一个 8KB 的缓冲数组那光池子本身就要占近 80MB 内存。如果并发峰值其实只需要 10 个对象池子却保留了几百个空闲对象内存就白白浪费了。解决思路是给池子设上限并且要区分“活跃对象”和“空闲对象”。空闲对象超过阈值就释放底层资源或者干脆清掉让 GC 回收。还有一种做法是使用弱引用持有空闲对象内存紧张时 JVM 可以自动回收恢复时再按需重建。5.4 问题四代码可读性与调试成本骤增最后这一个是很多团队转型时最大的阻力。可变对象、覆写式逻辑、重复使用的容器这些东西让代码的“数据流”变得不再直观。接手的人需要额外理解“这个对象此刻是新的还是旧的”“这个 List 清了吗”“这段逻辑会不会在某个分支里忘归还”。我的经验是把复用逻辑全部收敛到一个独立的模块里不要散落在业务代码中。对外提供的接口保持传统风格——我提供 borrow() 返回一个看似全新的对象提供 recycle() 回收对象。使用方不关心内部实现是新建还是复用只管借和还。这样即使内部用了原地重用业务代码的可读性也不会被破坏。再结合团队 code review 环节里加一条检查凡是持有了从池中借出的对象就必须在代码路径的出口处归还。这个规则写在文档里比临时讨论要靠谱得多。尽早建立约定后面维护成本就能压下来。最后的两个建议这次优化做完之后我脑子里的“性能意识”发生了不小的转变。以前写流式处理代码我想的是怎么把算法写得漂亮、把逻辑写清楚现在我会多问一句这段逻辑每触发一次堆上会多出几个对象它们活多久这个问题一问出来很多性能优化的答案就自己浮出水面了。最后再分享一个小技巧判断一段链路值不值得做原地重用不需要跑完整的压力测试直接抓一段堆分配的采样数据就行。Java 里可以用-XX:PrintGCDetails看 GC 频率或者用jcmd pid GC.class_histogram看各类型对象实例数。如果是别的运行时也有类似的 profiler 工具。看到某类对象实例数非常高、却几乎没有长生命周期实例——那就说明这里八成是“堆数据制造机”值得动手优化。原地重用真正考验的不是写代码的能力而是对数据生命周期的理解。生产链路上每一个对象从哪来、到哪去、什么时候能放回去心里有数代码才不会失控。希望这篇文章能帮你避开我踩过的那些坑让流式处理链路真正顺畅起来。
上一篇/下一篇内容由系统自动关联 返回资讯列表 →