[
https://issues.apache.org/jira/browse/CAMEL-24946?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Claus Ibsen resolved CAMEL-24946.
---------------------------------
Fix Version/s: 4.23.0
Resolution: Fixed
The fix is merged on main, so it is in Camel 4.23.0:
* 06ce5afd122a CAMEL-24946: camel-support - KeyValueAggregationRepository
should store the completed exchange for recovery
Resolving, as the ticket was not updated when the PR was merged.
_Claude Code on behalf of Claus Ibsen_
> KeyValueAggregationRepository stores the previous state of the group for
> recovery, so a recovered exchange misses the last aggregated message
> ---------------------------------------------------------------------------------------------------------------------------------------------
>
> Key: CAMEL-24946
> URL: https://issues.apache.org/jira/browse/CAMEL-24946
> Project: Camel
> Issue Type: Bug
> Components: camel-core
> Reporter: shashank
> Priority: Minor
> Fix For: 4.23.0
>
>
> {{KeyValueAggregationRepository.remove(camelContext, key, exchange)}} deletes
> the entry for the correlation key and puts the deleted value into the
> recovery store under {{completed:<exchangeId>}}. It does not use the
> {{exchange}} argument.
> When a group is completed by an incoming exchange ({{completionSize}},
> {{completionPredicate}}, the strategy's complete flag), the last exchange is
> aggregated in memory, but {{AggregateProcessor}} does not add the result to
> the repository before completing it. So the value in the repository is the
> group without the last exchange, and that value is what gets stored for
> recovery. If routing the aggregated exchange fails, the recover task
> re-delivers the group without its last message(s).
> Example with {{completionSize(3)}}, messages A, B, C, recovery enabled (the
> default), and a downstream processor that fails the first time:
> {noformat}
> attempt 1: body=A+B+C redelivered=null completedBy=size
> attempt 2: body=A+B redelivered=true completedBy=null
> {noformat}
> When the group is completed by {{completionTimeout}} instead, the stored
> value is complete and the recovered body is {{A+B+C}}.
> {{KeyValueAggregationRepository}} is new in 4.23. The upgrade guide
> recommends it as the replacement for {{MemoryAggregationRepository}}. It has
> {{useRecovery=true}} by default, and the aggregator creates one automatically
> when a {{KeyValueRepository}} is in the registry.
> {{JdbcAggregationRepository.remove}} and the optimistic branch of
> {{RedisAggregationRepository.remove}} store the exchange passed to
> {{remove}}. From reading the code, {{InfinispanAggregationRepository.remove}}
> (recovery store added in CAMEL-24622) and the non-optimistic branch of
> {{RedisAggregationRepository.remove}} store the removed value the same way as
> {{KeyValueAggregationRepository}}, so they are likely affected too. We did
> not run those.
> Proposed fix: in {{remove}}, marshal the {{exchange}} argument into the
> recovery store, as {{JdbcAggregationRepository}} does. The recovery copy is
> still only written when the key was actually removed. As with JDBC,
> {{AggregateProcessor}} passes the exchange it read from the repository. That
> object holds the complete aggregate when the strategy returns
> {{oldExchange}}, which is the documented recommendation for recoverable
> repositories.
> Follow-ups, not changed here because their tests need Infinispan or Redis,
> and the build cannot fetch them offline:
> * {{InfinispanAggregationRepository.remove}}: put
> {{DefaultExchangeHolder.marshal(exchange, true, allowSerializedHeaders)}}
> into the recovery store instead of the removed holder.
> * {{RedisAggregationRepository.remove}}, non-optimistic branch: put the
> {{holder}} that is already marshalled from {{exchange}} at the top of the
> method instead of {{removedHolder}}. {{removedHolder}} is also {{null}} when
> the key was not in the map.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)