[
https://issues.apache.org/jira/browse/FLINK-40326?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Aleksandr Savonin updated FLINK-40326:
--------------------------------------
Description:
GatedRateLimiter.acquire(int numberOfEvents) closes its gate only once
capacityLeft has reached zero:
{code:java}
if (capacityLeft <= 0) {
gatingFuture = new CompletableFuture<>();
return gatingFuture.thenRun(() -> capacityLeft -= numberOfEvents);
}
{code}
When 0 < capacityLeft < numberOfEvents the request is granted anyway and
capacityLeft goes negative, so the cycle admits more events than
capacityPerCycle.
With capacityPerCycle = 5 and 3 events per request, two requests admit 6. It
recurs every cycle, because notifyCheckpointComplete resets the capacity
instead of carrying the deficit.
The condition was correct while every request was for a single event. FLIP-535
changed the decrement from 1 to numberOfEvents, which is what turned the
per-cycle cap soft. Two smaller defects in the same class:
* acquire(int) does not validate numberOfEvents: a negative value increases
the remaining capacity, and zero consumes none. Not reachable from current
callers.
* notifyCheckpointComplete() dereferences the lazily initialised gatingFuture,
so a checkpoint completing before the first acquire() throws a
NullPointerException.
was:
GatedRateLimiter.acquire(int numberOfEvents) closes its gate only once
capacityLeft has reached zero:
{code:java}
if (capacityLeft <= 0) { gatingFuture = new CompletableFuture<>(); return
gatingFuture.thenRun(() -> capacityLeft -= numberOfEvents);
{code}
When 0 < capacityLeft < numberOfEvents the request is granted anyway and
capacityLeft goes negative, so the cycle admits more events than
capacityPerCycle.
With capacityPerCycle = 5 and 3 events per request, two requests admit 6. It
recurs every cycle, because notifyCheckpointComplete resets the capacity
instead of carrying the deficit.
The condition was correct while every request was for a single event. FLIP-535
changed the decrement from 1 to numberOfEvents, which is what turned the
per-cycle cap soft. Two smaller defects in the same class:
* acquire(int) does not validate numberOfEvents: a negative value increases the
remaining capacity, and zero consumes none. Not reachable from current callers.
* notifyCheckpointComplete() dereferences the lazily initialised gatingFuture,
so a checkpoint completing before the first acquire() throws a
NullPointerException.
> GatedRateLimiter admits more events per cycle than configured
> -------------------------------------------------------------
>
> Key: FLINK-40326
> URL: https://issues.apache.org/jira/browse/FLINK-40326
> Project: Flink
> Issue Type: Bug
> Components: Connectors / Common
> Reporter: Aleksandr Savonin
> Assignee: Aleksandr Savonin
> Priority: Major
>
> GatedRateLimiter.acquire(int numberOfEvents) closes its gate only once
> capacityLeft has reached zero:
> {code:java}
> if (capacityLeft <= 0) {
> gatingFuture = new CompletableFuture<>();
> return gatingFuture.thenRun(() -> capacityLeft -= numberOfEvents);
> }
> {code}
> When 0 < capacityLeft < numberOfEvents the request is granted anyway and
> capacityLeft goes negative, so the cycle admits more events than
> capacityPerCycle.
> With capacityPerCycle = 5 and 3 events per request, two requests admit 6. It
> recurs every cycle, because notifyCheckpointComplete resets the capacity
> instead of carrying the deficit.
> The condition was correct while every request was for a single event.
> FLIP-535 changed the decrement from 1 to numberOfEvents, which is what turned
> the per-cycle cap soft. Two smaller defects in the same class:
> * acquire(int) does not validate numberOfEvents: a negative value increases
> the remaining capacity, and zero consumes none. Not reachable from current
> callers.
> * notifyCheckpointComplete() dereferences the lazily initialised
> gatingFuture, so a checkpoint completing before the first acquire() throws a
> NullPointerException.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)