apache / apache/bifromq

heap memory leak because the ConcurrentLinkedQueue in Batcher is unbounded and filled

Open
#181 2 comments 0 reactions 0 assignees View on GitHub
Dominant language
Java
Stars
792
Forks
95
Avg merge
10h 45m
Merged PRs (30d)
4

Description

## ⚠️ *heap memory leak because the ConcurrentLinkedQueue in Batcher is unbounded and filled*

### **Describe the bug**

Image

Image

Image
the batcher_call_sched metric is normal, but batcher_call_time metric is zero, because the dist req is pushed into the callTaskBuffers, but in batchAndEmit() not polled. I suppose the grpc server did not response for a long time, then pipelineDepth.getAndDecrement() do not execute, so the pipelineDepth.get() < maxPipelineDepth.get() condition in trigger() can not be satisfied, so the dist req is pushed and not polled. then the heap is filled by the dist req.

```
public final CompletableFuture submit(BatcherKeyT batcherKey, CallT request) {
if (avgLatencyNanos.estimate() < burstLatencyNanos) {
ICallTask callTask = new CallTask<>(batcherKey, request);
boolean offered = callTaskBuffers.offer(callTask);
assert offered;
trigger();
return callTask.resultPromise();
} else {
dropCounter.increment();
return CompletableFuture.failedFuture(new BackPressureException("Too high average latency"));
}
}

private void trigger() {
if (triggering.compareAndSet(false, true)) {
try {
if (!callTaskBuffers.isEmpty() && pipelineDepth.get() < maxPipelineDepth.get()) {
batchAndEmit();
}
} catch (Throwable e) {
log.error("Unexpected exception", e);
} finally {
triggering.set(false);
if (!callTaskBuffers.isEmpty() && pipelineDepth.get() < maxPipelineDepth.get()) {
trigger();
}
}
}
}

private void batchAndEmit() {
pipelineDepth.incrementAndGet();
long buildStart = System.nanoTime();
IBatchCall batchCall = batchPool.poll();
assert batchCall != null;
int batchSize = 0;
LinkedList> batchedTasks = new LinkedList<>();
ICallTask callTask;
while (batchSize < maxBatchSize && (callTask = callTaskBuffers.poll()) != null) {
batchCall.add(callTask);
batchedTasks.add(callTask);
batchSize++;
queueingTimeSummary.record(System.nanoTime() - callTask.ts());
}
batchSizeSummary.record(batchSize);
long execStart = System.nanoTime();
batchBuildTimeSummary.record((execStart - buildStart));
final int finalBatchSize = batchSize;
batchCall.execute()
.whenComplete((v, e) -> {
long execEnd = System.nanoTime();
if (e != null) {
log.error("Unexpected exception during handling batchcall result", e);
// reset max batch size
maxBatchSize = 1;
} else {
long thisLatency = execEnd - execStart;
if (thisLatency > 0) {
updateMaxBatchSize(finalBatchSize, thisLatency);
}
batchExecTimer.record(thisLatency, TimeUnit.NANOSECONDS);
}
batchedTasks.forEach(t -> {
long callLatency = execEnd - t.ts();
avgLatencyNanos.observe(callLatency);
batchCallTimer.record(callLatency, TimeUnit.NANOSECONDS);
});
batchCall.reset();
batchPool.offer(batchCall);
pipelineDepth.getAndDecrement();
if (!callTaskBuffers.isEmpty()) {
trigger();
}
});
}
```

#### **Environment**

- Version: [3.2.1]
- JVM Version: [OpenJDK17]
- Hardware Spec: [32c64g]
- OS: [CentOS 7]

Contributor guide

No contributing guide indexed for this repository

Research direction

Start by tracing Batcher's submit(), trigger(), and batchAndEmit() methods, focusing on how pipelineDepth changes when a gRPC request is slow to complete. The issue names no test or file path; done means addressing the reported unbounded accumulation of queued requests when batching is blocked.

Written by the indexing model from the issue text.

Assessment

Tech stack
grpc, java
Domain
backend, distributed-systems
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.