[
https://issues.apache.org/jira/browse/CAMEL-24945?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Claus Ibsen updated CAMEL-24945:
--------------------------------
Fix Version/s: 4.23.0
> Aggregate EIP: with optimistic locking a group can lose its completion
> timeout and never complete
> -------------------------------------------------------------------------------------------------
>
> Key: CAMEL-24945
> URL: https://issues.apache.org/jira/browse/CAMEL-24945
> Project: Camel
> Issue Type: Bug
> Components: camel-core
> Reporter: shashank
> Priority: Minor
> Fix For: 4.23.0
>
>
> With {{optimisticLocking()}} and {{completionTimeout}}, a new group can be
> left in the repository without a timeout. It then only completes if another
> exchange arrives for the same key, or when the route is stopped with
> {{forceCompletionOnStop}}. With an in-memory repository it is lost when the
> route stops.
> The timeout map is keyed by correlation key. When a group is completed (for
> example by {{completionSize}}), {{onCompletion}} first removes the group from
> the repository and then calls {{timeoutMap.remove(key)}}. With optimistic
> locking there is no lock around these two calls. If another exchange for the
> same key starts a new group in between, it registers its timeout
> ({{trackTimeout}} runs before the repository add) and adds the group. The
> completing thread then removes that new timeout entry.
> Example with {{completionSize(2).completionTimeout(500)}}, key 1:
> # {{a}} is aggregated.
> # {{b}} completes [a, b] by size. Its thread removes the group from the
> repository.
> # {{c}} arrives on another thread, registers a timeout for key 1 and adds
> group [c].
> # The thread of {{b}} calls {{timeoutMap.remove("1")}}, which removes the
> timeout of [c].
> Observed with {{MemoryAggregationRepository(true)}} and the thread of {{b}}
> paused after its {{remove()}}, waiting 2.5 seconds afterwards:
> {noformat}
> repo before final flush={1=c}
> downstream output: [a+b (completedBy=size)]
> {noformat}
> Without the pause: {{[a+b (completedBy=size), c (completedBy=timeout)]}}.
> In pessimistic mode both calls are made under the aggregation lock, so this
> only affects optimistic locking.
> Proposed fix: with optimistic locking, do not remove the timeout entry in
> {{onCompletion}}. A leftover entry is harmless because the timeout eviction
> reads the group and removes it with a compare-and-set, and every new group
> refreshes the entry before it is added. In the same mode, {{onEviction}}
> should not skip an eviction because the exchange id in the entry is in
> progress: with leftover entries that id can belong to a completed group while
> a newer group for the key is waiting. The compare-and-set already prevents
> completing a group twice. We checked this change in a TLA+ model. Two other
> options still leave a group without a timeout: removing the entry only when
> it holds an exchange id of the completed group, and registering the timeout
> after the add. The second part is needed too: the model, and a unit test with
> a repository that keeps exchange ids (like JDBC), strand a group when only
> the first part is applied. For example, m1 and m2 start a new group at the
> same time, and m2's timeout entry overwrites m1's. m3 then completes [m2,
> m3], and m1's add succeeds afterwards. The entry for the key still holds m2's
> id, which is in progress, so the eviction is skipped and [m1] never times out.
> Side effects, with optimistic locking only: after a group completes, its
> timeout entry stays until it expires. The eviction then does one repository
> {{get}} for the key, which returns nothing or a newer group. A newer group
> created on the same Camel instance always refreshes the entry. With several
> instances sharing the repository, an old entry on one instance can complete a
> group that another instance created for the same key within
> {{completionTimeout}}, so it completes up to one timeout early.
> Cross-instance timeouts already work this way for exchanges added to an
> existing group from another instance.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)