apache / apache/pulsar-client-go

[Bug] Pulsar Go Client consumer memory limit does not work as expected

Open
#1,427 2 comments 0 reactions 0 assignees View on GitHub
Dominant language
Go
Stars
745
Forks
389
Avg merge
3d 20h
Merged PRs (30d)
3

Description

#### Expected behavior

When the consumer receives messages, the memory reserved by `memLimit.ForceReserveMemory` should be fully released by `memLimit.ReleaseMemory`, ensuring the memory counter remains balanced and does not exceed the configured limit.

#### Actual behavior

In `consumer_partition.go`, the logic is inconsistent:

* At [consumer_partition.go#L1181-L1436](https://github.com/apache/pulsar-client-go/blob/v0.16.0/pulsar/consumer_partition.go#L1181-L1436), `memLimit.ForceReserveMemory` uses the accumulated size of all messages.
* At [consumer_partition.go#L1587-L1708](https://github.com/apache/pulsar-client-go/blob/v0.16.0/pulsar/consumer_partition.go#L1587-L1708), `memLimit.ReleaseMemory` is called with only the raw/original size.

This leads to a mismatch. For example, if a batch of 3 messages (each 1 KiB) is received:

* `memLimit.ForceReserveMemory` reserves **6 KiB**.
* Later, `memLimit.ReleaseMemory` releases only **3 KiB**.

As a result, the memory usage counter can grow beyond the actual limit.

#### Steps to reproduce

1. Run a Pulsar Go consumer with `memLimit` configured.
2. Consume batched messages where the aggregated size differs from the raw size (e.g., multiple small messages).
3. Observe that memory usage grows inconsistently and eventually exceeds the expected limit.

#### System configuration

**Pulsar Go Client version**: v0.14.0, v0.15.1, v0.16.0 (and possibly others)

#### Suggested fix

Ensure that `memLimit.ReleaseMemory` is called with the same size used by `memLimit.ForceReserveMemory`.

#### Relevant code

- [consumer_partition.go#L1181-L1436](https://github.com/apache/pulsar-client-go/blob/v0.16.0/pulsar/consumer_partition.go#L1181-L1436)

```go
func (pc *partitionConsumer) MessageReceived(response *pb.CommandMessage, headersAndPayload internal.Buffer) error {
...
var (
bytesReceived int
skippedMessages int32
)
for i := 0; i < numMsgs; i++ {
...
messages = append(messages, msg)
bytesReceived += msg.size()
if pc.options.autoReceiverQueueSize {
pc.client.memLimit.ForceReserveMemory(int64(bytesReceived)) // should be int64(msg.size())
pc.incomingMessages.Add(int32(1))
pc.markScaleIfNeed()
}
}
...
}
```

- [consumer_partition.go#L1587-L1708](https://github.com/apache/pulsar-client-go/blob/v0.16.0/pulsar/consumer_partition.go#L1587-L1708)

```go
func (pc *partitionConsumer) dispatcher() {
...
for {
...
select {
...
// if the messageCh is nil or the messageCh is full this will not be selected
case messageCh <- nextMessage:
// allow this message to be garbage collected
messages[0] = nil
messages = messages[1:]

// for the zeroQueueConsumer, the permits controlled by itself
if pc.options.receiverQueueSize > 0 {
pc.availablePermits.inc()
}

if pc.options.autoReceiverQueueSize {
pc.incomingMessages.Dec()
pc.client.memLimit.ReleaseMemory(int64(nextMessageSize)) // this is right
pc.expectMoreIncomingMessages()
}
...
}
...
}
...
}
```

Contributor guide

Open the contributing guide

Research direction

Start in pulsar/consumer_partition.go at MessageReceived around lines 1181-1436, then trace dispatcher around lines 1587-1708. Compare the size passed to memLimit.ForceReserveMemory with the size passed to ReleaseMemory during batched delivery. Done means the accounting uses matching sizes so the memory counter remains balanced under the reproduction described.

Written by the indexing model from the issue text.

Assessment

Tech stack
go
Domain
distributed-systems
Issue type
Bug
Difficulty
2/5
Estimated time
1-3 hours
Activity status
Stale
Clarity
Clearly specified
Newbie friendliness
52/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.