This is an automated email from the ASF dual-hosted git repository.
oscerd 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 72e9fb578aa0 CAMEL-24493: camel-ibm-cos - honor the multiPartUpload,
partSize, storageClass and deleteAfterWrite producer options (#25750)
72e9fb578aa0 is described below
commit 72e9fb578aa0c9f31fc56a6e9a6e4727ab2a0b29
Author: Andrea Cosentino <[email protected]>
AuthorDate: Wed Aug 26 09:14:38 2026 +0200
CAMEL-24493: camel-ibm-cos - honor the multiPartUpload, partSize,
storageClass and deleteAfterWrite producer options (#25750)
IBMCOSProducer.putObject() built a plain PutObjectRequest and ignored the
four
producer options declared on IBMCOSConfiguration and advertised in the docs.
Apply them, mirroring camel-aws2-s3: set storageClass on the request,
upload via
the SDK TransferManager (minimum part size / multipart threshold =
partSize) when
multiPartUpload is enabled, and delete the local File payload after a
successful
upload when deleteAfterWrite is set.
Signed-off-by: Andrea Cosentino <[email protected]>
Co-authored-by: Claude Opus 4.8 <[email protected]>
---
.../camel/component/ibm/cos/IBMCOSProducer.java | 51 ++++++++++++++++++++--
1 file changed, 48 insertions(+), 3 deletions(-)
diff --git
a/components/camel-ibm/camel-ibm-cos/src/main/java/org/apache/camel/component/ibm/cos/IBMCOSProducer.java
b/components/camel-ibm/camel-ibm-cos/src/main/java/org/apache/camel/component/ibm/cos/IBMCOSProducer.java
index e8bed18528fa..c4e105af2e00 100644
---
a/components/camel-ibm/camel-ibm-cos/src/main/java/org/apache/camel/component/ibm/cos/IBMCOSProducer.java
+++
b/components/camel-ibm/camel-ibm-cos/src/main/java/org/apache/camel/component/ibm/cos/IBMCOSProducer.java
@@ -16,6 +16,7 @@
*/
package org.apache.camel.component.ibm.cos;
+import java.io.File;
import java.io.InputStream;
import java.util.ArrayList;
import java.util.List;
@@ -33,10 +34,14 @@ import
com.ibm.cloud.objectstorage.services.s3.model.ObjectMetadata;
import com.ibm.cloud.objectstorage.services.s3.model.PutObjectRequest;
import com.ibm.cloud.objectstorage.services.s3.model.PutObjectResult;
import com.ibm.cloud.objectstorage.services.s3.model.S3Object;
+import com.ibm.cloud.objectstorage.services.s3.transfer.TransferManager;
+import com.ibm.cloud.objectstorage.services.s3.transfer.TransferManagerBuilder;
+import com.ibm.cloud.objectstorage.services.s3.transfer.model.UploadResult;
import org.apache.camel.Exchange;
import org.apache.camel.Message;
import org.apache.camel.WrappedFile;
import org.apache.camel.support.DefaultProducer;
+import org.apache.camel.util.FileUtil;
import org.apache.camel.util.ObjectHelper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -154,12 +159,52 @@ public class IBMCOSProducer extends DefaultProducer {
PutObjectRequest putObjectRequest = new PutObjectRequest(bucketName,
key, inputStream, metadata);
+ String storageClass = getConfiguration().getStorageClass();
+ if (ObjectHelper.isNotEmpty(storageClass)) {
+ putObjectRequest.withStorageClass(storageClass);
+ }
+
LOG.trace("Putting object [{}] into bucket [{}]...", key, bucketName);
- PutObjectResult putObjectResult =
cosClient.putObject(putObjectRequest);
+
+ String eTag;
+ String versionId;
+ if (getConfiguration().isMultiPartUpload()) {
+ // Upload via the TransferManager so large payloads are sent as a
multipart upload with the configured part size
+ TransferManager transferManager = TransferManagerBuilder.standard()
+ .withS3Client(cosClient)
+
.withMinimumUploadPartSize(getConfiguration().getPartSize())
+
.withMultipartUploadThreshold(getConfiguration().getPartSize())
+ .build();
+ try {
+ UploadResult uploadResult =
transferManager.upload(putObjectRequest).waitForUploadResult();
+ eTag = uploadResult.getETag();
+ versionId = uploadResult.getVersionId();
+ } finally {
+ // shut down the TransferManager's own thread pool but keep
the shared COS client open
+ transferManager.shutdownNow(false);
+ }
+ } else {
+ PutObjectResult putObjectResult =
cosClient.putObject(putObjectRequest);
+ eTag = putObjectResult.getETag();
+ versionId = putObjectResult.getVersionId();
+ }
Message message = getMessageForResponse(exchange);
- message.setHeader(IBMCOSConstants.E_TAG, putObjectResult.getETag());
- message.setHeader(IBMCOSConstants.VERSION_ID,
putObjectResult.getVersionId());
+ message.setHeader(IBMCOSConstants.E_TAG, eTag);
+ message.setHeader(IBMCOSConstants.VERSION_ID, versionId);
+
+ if (getConfiguration().isDeleteAfterWrite()) {
+ File filePayload = null;
+ if (body instanceof File file) {
+ filePayload = file;
+ } else if (body instanceof WrappedFile<?> wrapped &&
wrapped.getFile() instanceof File file) {
+ filePayload = file;
+ }
+ if (filePayload != null) {
+ LOG.trace("Deleting file payload [{}] after write",
filePayload);
+ FileUtil.deleteFile(filePayload);
+ }
+ }
}
private void getObject(AmazonS3 cosClient, Exchange exchange) {