This is an automated email from the ASF dual-hosted git repository.
davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new d8be2daaa324 CAMEL-24195: camel-aws2-s3 - Fix StreamUploadProducer
race condition and timeout task exception handling
d8be2daaa324 is described below
commit d8be2daaa32416f2941282f5fc9aa9dad4598cfe
Author: Claus Ibsen <[email protected]>
AuthorDate: Fri Jul 17 22:42:12 2026 +0200
CAMEL-24195: camel-aws2-s3 - Fix StreamUploadProducer race condition and
timeout task exception handling
Co-Authored-By: Claude Opus 4.6 <[email protected]>
---
.../aws2/s3/stream/AWS2S3StreamUploadProducer.java | 29 ++++++++++++++++------
1 file changed, 21 insertions(+), 8 deletions(-)
diff --git
a/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/stream/AWS2S3StreamUploadProducer.java
b/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/stream/AWS2S3StreamUploadProducer.java
index b75ba975c060..0ed57fd2e99c 100644
---
a/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/stream/AWS2S3StreamUploadProducer.java
+++
b/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/stream/AWS2S3StreamUploadProducer.java
@@ -137,8 +137,12 @@ public class AWS2S3StreamUploadProducer extends
DefaultProducer {
lock.lock();
try {
if (ObjectHelper.isNotEmpty(uploadAggregate)) {
- uploadPart(uploadAggregate);
- completeUpload(uploadAggregate);
+ try {
+ uploadPart(uploadAggregate);
+ completeUpload(uploadAggregate);
+ } catch (Exception e) {
+ LOG.warn("Error during timeout-triggered upload flush
for {}", uploadAggregate.dynamicKeyName, e);
+ }
uploadAggregate = null;
}
@@ -148,8 +152,12 @@ public class AWS2S3StreamUploadProducer extends
DefaultProducer {
for (Map.Entry<Long, UploadState> entry :
timestampBasedUploads.entrySet()) {
UploadState state = entry.getValue();
if (ObjectHelper.isNotEmpty(state) &&
state.buffer.size() > 0) {
- uploadPart(state);
- completeUpload(state);
+ try {
+ uploadPart(state);
+ completeUpload(state);
+ } catch (Exception e) {
+ LOG.warn("Error during timeout-triggered
upload flush for {}", state.dynamicKeyName, e);
+ }
keysToRemove.add(entry.getKey());
}
}
@@ -178,10 +186,15 @@ public class AWS2S3StreamUploadProducer extends
DefaultProducer {
byte[] b;
int maxRead = (getConfiguration().isMultiPartUpload()
? Math.toIntExact(getConfiguration().getPartSize()) :
getConfiguration().getBufferSize());
- if (uploadAggregate != null) {
- uploadAggregate.index++;
- maxRead -= uploadAggregate.buffer.size();
- maxRead = Math.max(1, maxRead);
+ lock.lock();
+ try {
+ if (uploadAggregate != null) {
+ uploadAggregate.index++;
+ maxRead -= uploadAggregate.buffer.size();
+ maxRead = Math.max(1, maxRead);
+ }
+ } finally {
+ lock.unlock();
}
while ((b = AWS2S3Utils.toByteArray(is, maxRead)) != null && b.length
> 0) {