[
https://issues.apache.org/jira/browse/CAMEL-25130?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Claus Ibsen resolved CAMEL-25130.
---------------------------------
Resolution: Fixed
Fixed via https://github.com/apache/camel/pull/27039 (merged to main for
4.23.0).
_Claude Code on behalf of davsclaus_
> camel-sql - JdbcAggregationRepository with optimistic locking: add() ignores
> the old exchange, so a group completed meanwhile is stored again (duplicates)
> and a new group can be overwritten (lost messages)
> -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
>
> Key: CAMEL-25130
> URL: https://issues.apache.org/jira/browse/CAMEL-25130
> Project: Camel
> Issue Type: Bug
> Components: camel-sql
> Reporter: shashank
> Priority: Minor
> Fix For: 4.23.0
>
>
> The Aggregate EIP with {{optimisticLocking()}} and
> {{JdbcAggregationRepository}} is the documented set-up for several Camel
> instances that share the aggregation table ("Optimistic locking" in the
> camel-sql documentation). {{AggregateProcessor.doAggregation}} reads the
> group with {{repository.get}} (which puts the row version in the
> {{CamelOptimisticLockVersion}} property), aggregates, and then calls
> {{OptimisticLockingAggregationRepository.add(camelContext, key, oldExchange,
> newExchange)}}. The contract of that method is a compare-and-set against
> {{oldExchange}}: when {{oldExchange}} is null there must be no stored
> exchange, otherwise the stored exchange must still be the one that was read.
> With optimistic locking the aggregator does not lock the key, so another
> thread or node can complete the group between {{get}} and {{add}}.
> {{JdbcAggregationRepository.add(camelContext, key, oldExchange,
> newExchange)}} ignores {{oldExchange}} and calls {{add(camelContext, key,
> newExchange)}}, which does in one transaction:
> * {{SELECT COUNT(1) ... WHERE id = ?}};
> * if a row exists: {{UPDATE ... WHERE id = ? AND version = ?}} with the
> version of the new exchange (0 rows: {{OptimisticLockingException}});
> * if no row exists: {{INSERT}} with version 1.
> Two cases are not detected:
> *1. The group was completed while this thread aggregated on it.* Thread A
> reads group [x] and aggregates its message a. Meanwhile the group completes
> on another thread or node: a message that satisfies the completion predicate
> or size, or the completion timeout checker. {{remove()}} deletes the row (the
> version matches) and [x, b] is sent. Thread A then calls {{add}}: the row is
> gone, so it inserts [x, a] as a new group. The next completion sends [x, a,
> ...]: x is delivered twice. Every message of the completed group is sent
> again with the next group.
> *2. Every group starts at version 1 (ABA).* As above, but a message c starts
> a new group for the key (INSERT, version 1) before thread A calls {{add}}.
> Thread A's {{UPDATE ... WHERE version = 1}} matches the new group's row and
> replaces [c] with [x, a]: c is lost without an error, and x is delivered
> twice. {{remove()}} has the same problem: a completion that read version 1 of
> the old group can delete the new group.
> The window is the time between {{get}} and {{add}} (the aggregation strategy
> plus two queries). It is hit when messages for a key keep arriving while its
> group completes, which is normal under load.
> h3. Reproduction
> Embedded H2, a {{DataSourceTransactionManager}}, the real
> {{JdbcAggregationRepository}} (tables as in the documentation), and a route
> {{from("direct:in").aggregate(header("id"),
> concat).aggregationRepository(repo).optimisticLocking().eagerCheckCompletion().completionPredicate(header("last").isEqualTo("true"))}}.
> The aggregation strategy returns the old exchange (the recommended style)
> and holds message "a" on a latch after {{get}} has read the group; no sleeps
> are needed for the ordering.
> * completed by another message: send x; thread A sends a (held); send b with
> last=true (completes [x,b]); release A; send z with last=true. Sent groups:
> [x,b] and [x,a,z]. x is delivered twice, 3 of 3 runs.
> * completed by the timeout: {{completionTimeout(300)}}, send x, thread A
> sends a (held), the timeout sends [x], release A. Sent groups: [x] and [x,a],
> 3 of 3 runs.
> * ABA: as the first case, plus send c (a new group) before releasing A. Sent
> groups: [x,b] and [x,a,z]. c is lost and x is delivered twice, 3 of 3 runs.
> * Control: the same messages without overlap: [x,a,b], [z].
> A TLA+ model of the repository (each repository call one atomic transaction,
> which is stronger than READ COMMITTED) with aggregating threads, a completing
> message and the timeout checker checks that no message is sent twice and none
> is lost. It finds case 1 in 4 steps (Get(a), Get(b), Complete(b), Add(a)) and
> the timeout variant, and case 2 when only the check for a removed row is
> added. With both parts of the fix below the properties hold for three threads
> plus the timeout checker, also starting from an empty table.
> h3. Proposed fix
> Make {{add(camelContext, key, oldExchange, newExchange)}} a compare-and-set
> against {{oldExchange}}:
> * the expected version is the one of {{oldExchange}} (its
> {{CamelOptimisticLockVersion}}), not the one of {{newExchange}};
> * {{oldExchange == null}} and a row exists: {{OptimisticLockingException}}
> (as today);
> * {{oldExchange != null}} and no row: {{OptimisticLockingException}} (new),
> so the aggregator retries and starts a new group with only its own message;
> * a new group starts with a random positive version instead of 1, so a
> version read from an earlier group of the same key never matches the new
> group. This also stops a stale {{remove()}} from deleting a new group.
> The schema is unchanged: the documented column is {{version BIGINT NOT
> NULL}}. Groups stored before the upgrade keep their version. The
> non-optimistic {{add(camelContext, key, exchange)}} keeps its behaviour (the
> aggregator's lock serializes it); only the version stored for a new group
> changes. {{ClusteredJdbcAggregationRepository}},
> {{PostgresAggregationRepository}} and
> {{ClusteredPostgresAggregationRepository}} inherit {{add}}. Recovery does not
> compare versions and is not affected. A side effect of taking the version
> from the old exchange: an aggregation strategy that returns the new exchange
> no longer fails every optimistic {{add}}.
> Tests: the three cases above as unit tests with a latch in the aggregation
> strategy, and repository-level tests for add after a removal, add after a new
> group was started, and a stale remove of a new group. A note for the 4.23
> upgrade guide (new groups no longer start at version 1).
> Affected: all versions with optimistic locking for this repository (the same
> {{add}} code is at the camel-3.4.0, 4.0.0, 4.14.0 and 4.18.0 tags). Before
> 4.22.0 {{remove()}} also ignored its delete count (CAMEL-24048).
> Duplicate check (2026-09-29): JIRA text "JdbcAggregationRepository" (31
> issues, including CAMEL-24048 stale remove, CAMEL-17613 confirm order,
> CAMEL-7810 and CAMEL-6144 optimistic locking), "aggregation repository" with
> "version", and component camel-sql since June 2026: none covers add() after a
> completion or the version restarting at 1. GitHub pull requests for
> "JdbcAggregationRepository" and "aggregation repository optimistic":
> apache/camel#24710 (CAMEL-24048), #24655, #24707, #26791, not this.
> _Filed with Claude Code on behalf of allthingssecurity._
--
This message was sent by Atlassian Jira
(v8.20.10#820010)