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 17c0630c108a CAMEL-25167: Large payloads - consistent body length
detection and streaming charset conversion in file / object-storage producers
(#27132)
17c0630c108a is described below
commit 17c0630c108a3018fc8a02e093710943e2f9e32d
Author: Federico Mariani <[email protected]>
AuthorDate: Thu Oct 1 07:52:48 2026 +0200
CAMEL-25167: Large payloads - consistent body length detection and
streaming charset conversion in file / object-storage producers (#27132)
- camel-support: add PayloadHelper to determine the length of a body
without reading it, and to copy a stream of unknown length through
stream caching (spooled to disk when spooling is enabled) instead of
a byte array in memory. Only the length of a stream cache that is a
byte stream is trusted (a ReaderCache length counts characters).
- aws2-s3, azure-storage-blob, azure-storage-datalake, minio, ibm-cos,
google-storage: use PayloadHelper instead of the mark/skip/reset
length probes, which failed with "Resetting to invalid mark" on
buffered streams (such as a java.nio.file.Path body), and instead of
reading unknown-length bodies into memory. Datalake now accepts
streams without mark support, and blob/datalake accept remote files.
AWS2S3Utils.determineLengthInputStream, BlobUtils.getInputStreamLength
and DataLakeUtils.getInputStreamLength are deprecated.
- camel-util: add IOHelper.ReaderInputStream, which encodes a Reader
chunk by chunk with a single CharsetEncoder (a byte order mark is
written once); EncodingInputStream now builds on it.
- ftp, sftp, smb: with charset configured, encode the body while it is
uploaded instead of loading it as a whole String.
- Add a "Large payloads" page to the user manual, linked from the
Stream caching page.
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../camel/component/aws2/s3/AWS2S3Producer.java | 41 +++--
.../camel/component/aws2/s3/utils/AWS2S3Utils.java | 40 +---
.../aws2/s3/AWS2S3ProducerPayloadLengthTest.java | 152 +++++++++++++++
.../component/azure/storage/blob/BlobBlock.java | 13 +-
.../azure/storage/blob/BlobStreamAndLength.java | 58 +++---
.../component/azure/storage/blob/BlobUtils.java | 40 +---
.../storage/blob/BlobStreamAndLengthTest.java | 52 ++++++
.../azure/storage/datalake/DataLakeUtils.java | 18 +-
.../storage/datalake/FileStreamAndLength.java | 52 +++---
.../storage/datalake/FileStreamAndLengthTest.java | 57 ++++++
.../camel/component/file/GenericFileHelper.java | 31 ++++
.../component/file/GenericFileHelperTest.java | 45 +++++
.../camel/component/file/remote/FtpOperations.java | 2 +-
.../component/file/remote/SftpOperations.java | 5 +-
.../integration/SftpProducerWithCharsetIT.java | 18 ++
.../google/storage/GoogleCloudStorageProducer.java | 47 +++--
.../camel/component/ibm/cos/IBMCOSProducer.java | 21 ++-
.../camel/component/minio/MinioProducer.java | 44 ++---
.../apache/camel/component/smb/SmbOperations.java | 2 +-
.../apache/camel/support/PayloadHelperTest.java | 124 +++++++++++++
.../org/apache/camel/support/PayloadHelper.java | 119 ++++++++++++
.../main/java/org/apache/camel/util/IOHelper.java | 126 +++++++++----
docs/user-manual/modules/ROOT/nav.adoc | 1 +
docs/user-manual/modules/ROOT/pages/index.adoc | 1 +
.../modules/ROOT/pages/large-payloads.adoc | 205 +++++++++++++++++++++
.../modules/ROOT/pages/stream-caching.adoc | 2 +
26 files changed, 1060 insertions(+), 256 deletions(-)
diff --git
a/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/AWS2S3Producer.java
b/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/AWS2S3Producer.java
index da44011edbe5..680832b5389a 100644
---
a/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/AWS2S3Producer.java
+++
b/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/AWS2S3Producer.java
@@ -16,8 +16,6 @@
*/
package org.apache.camel.component.aws2.s3;
-import java.io.ByteArrayInputStream;
-import java.io.ByteArrayOutputStream;
import java.io.File;
import java.io.FileInputStream;
import java.io.InputStream;
@@ -36,6 +34,7 @@ import org.apache.camel.Message;
import org.apache.camel.WrappedFile;
import org.apache.camel.component.aws2.s3.utils.AWS2S3Utils;
import org.apache.camel.support.DefaultProducer;
+import org.apache.camel.support.PayloadHelper;
import org.apache.camel.util.FileUtil;
import org.apache.camel.util.IOHelper;
import org.apache.camel.util.ObjectHelper;
@@ -192,18 +191,19 @@ public class AWS2S3Producer extends DefaultProducer {
contentLength = f.length();
} else {
// okay we use input stream
+ if (contentLength <= 0) {
+ // such as a java.nio.file.Path, byte[] or stream cache
+ contentLength = PayloadHelper.getLength(obj);
+ }
inputStream = exchange.getIn().getMandatoryBody(InputStream.class);
if (contentLength <= 0) {
- contentLength =
AWS2S3Utils.determineLengthInputStream(inputStream);
+ contentLength = PayloadHelper.getLength(inputStream);
if (contentLength == -1) {
- // fallback to read into memory to calculate length
- LOG.debug(
- "The content length is not defined. It needs to be
determined by reading the data into memory");
- ByteArrayOutputStream baos = new ByteArrayOutputStream();
- IOHelper.copyAndCloseInput(inputStream, baos);
- byte[] arr = baos.toByteArray();
- contentLength = arr.length;
- inputStream = new ByteArrayInputStream(arr);
+ // fallback to copy the data to calculate the length,
which uses stream caching
+ // so big payloads are spooled to disk when spooling is
enabled
+ LOG.debug("The content length is not defined. It needs to
be determined by copying the data");
+ inputStream = PayloadHelper.cacheStream(exchange,
inputStream);
+ contentLength = PayloadHelper.getLength(inputStream);
}
}
}
@@ -369,18 +369,19 @@ public class AWS2S3Producer extends DefaultProducer {
contentLength = filePayload.length();
} else {
// okay we use input stream
+ if (contentLength <= 0) {
+ // such as a java.nio.file.Path, byte[] or stream cache
+ contentLength = PayloadHelper.getLength(obj);
+ }
inputStream =
exchange.getIn().getMandatoryBody(InputStream.class);
if (contentLength <= 0) {
- contentLength =
AWS2S3Utils.determineLengthInputStream(inputStream);
+ contentLength = PayloadHelper.getLength(inputStream);
if (contentLength == -1) {
- // fallback to read into memory to calculate length
- LOG.debug(
- "The content length is not defined. It needs
to be determined by reading the data into memory");
- ByteArrayOutputStream baos = new
ByteArrayOutputStream();
- IOHelper.copyAndCloseInput(inputStream, baos);
- byte[] arr = baos.toByteArray();
- contentLength = arr.length;
- inputStream = new ByteArrayInputStream(arr);
+ // fallback to copy the data to calculate the length,
which uses stream caching
+ // so big payloads are spooled to disk when spooling
is enabled
+ LOG.debug("The content length is not defined. It needs
to be determined by copying the data");
+ inputStream = PayloadHelper.cacheStream(exchange,
inputStream);
+ contentLength = PayloadHelper.getLength(inputStream);
}
}
}
diff --git
a/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/utils/AWS2S3Utils.java
b/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/utils/AWS2S3Utils.java
index f6cb21e4457b..13c94bf8d366 100644
---
a/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/utils/AWS2S3Utils.java
+++
b/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/utils/AWS2S3Utils.java
@@ -16,17 +16,15 @@
*/
package org.apache.camel.component.aws2.s3.utils;
-import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
-import java.io.FileInputStream;
import java.io.IOException;
import java.io.InputStream;
import org.apache.camel.Exchange;
-import org.apache.camel.StreamCache;
import org.apache.camel.component.aws2.s3.AWS2S3Configuration;
import org.apache.camel.component.aws2.s3.AWS2S3Constants;
import org.apache.camel.spi.Language;
+import org.apache.camel.support.PayloadHelper;
import org.apache.camel.util.ObjectHelper;
import software.amazon.awssdk.services.s3.model.CreateMultipartUploadRequest;
import software.amazon.awssdk.services.s3.model.ServerSideEncryption;
@@ -95,35 +93,15 @@ public final class AWS2S3Utils {
}
}
+ /**
+ * Determines the length of the input stream, without reading it.
+ *
+ * @return the length, or <tt>-1</tt> if the length cannot be
determined without reading the stream
+ * @deprecated use {@link PayloadHelper#getLength(Object)}
+ */
+ @Deprecated(since = "4.23.0")
public static long determineLengthInputStream(InputStream is) throws
IOException {
- if (is instanceof StreamCache streamCache) {
- long len = streamCache.length();
- if (len > 0) {
- return len;
- }
- } else if (is instanceof FileInputStream fis) {
- return fis.getChannel().size();
- }
-
- if (!is.markSupported()) {
- return -1;
- }
- if (is instanceof ByteArrayInputStream) {
- return is.available();
- }
- long size = 0;
- try {
- is.mark(1024);
- int i = is.available();
- while (i > 0) {
- long skip = is.skip(i);
- size += skip;
- i = is.available();
- }
- } finally {
- is.reset();
- }
- return size;
+ return PayloadHelper.getLength(is);
}
public static byte[] toByteArray(InputStream is, final int size) throws
IOException {
diff --git
a/components/camel-aws/camel-aws2-s3/src/test/java/org/apache/camel/component/aws2/s3/AWS2S3ProducerPayloadLengthTest.java
b/components/camel-aws/camel-aws2-s3/src/test/java/org/apache/camel/component/aws2/s3/AWS2S3ProducerPayloadLengthTest.java
new file mode 100644
index 000000000000..581fc0bdb9a3
--- /dev/null
+++
b/components/camel-aws/camel-aws2-s3/src/test/java/org/apache/camel/component/aws2/s3/AWS2S3ProducerPayloadLengthTest.java
@@ -0,0 +1,152 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.component.aws2.s3;
+
+import java.io.ByteArrayInputStream;
+import java.io.FilterInputStream;
+import java.io.InputStream;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.Random;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.converter.stream.ReaderCache;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.camel.support.DefaultExchange;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+import org.mockito.Mock;
+import org.mockito.MockitoAnnotations;
+import software.amazon.awssdk.core.sync.RequestBody;
+import software.amazon.awssdk.http.SdkHttpResponse;
+import software.amazon.awssdk.services.s3.S3Client;
+import software.amazon.awssdk.services.s3.model.PutObjectRequest;
+import software.amazon.awssdk.services.s3.model.PutObjectResponse;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.when;
+
+/**
+ * Tests that the producer determines the length of a body that is not a
{@link java.io.File}, without failing on big
+ * payloads and without losing content.
+ */
+public class AWS2S3ProducerPayloadLengthTest {
+
+ // bigger than the buffer of a BufferedInputStream, which a mark/reset
based length probe cannot handle
+ private static final int PAYLOAD_SIZE = 1024 * 1024;
+
+ @Mock
+ private AWS2S3Endpoint endpoint;
+
+ @Mock
+ private AWS2S3Configuration configuration;
+
+ @Mock
+ private S3Client s3Client;
+
+ private AWS2S3Producer producer;
+ private DefaultCamelContext camelContext;
+ private byte[] payload;
+
+ private Long uploadedLength;
+ private byte[] uploadedContent;
+
+ @BeforeEach
+ public void setUp() throws Exception {
+ MockitoAnnotations.openMocks(this);
+ camelContext = new DefaultCamelContext();
+
+ when(endpoint.getConfiguration()).thenReturn(configuration);
+ when(endpoint.getCamelContext()).thenReturn(camelContext);
+ when(endpoint.getS3Client()).thenReturn(s3Client);
+ when(configuration.getBucketName()).thenReturn("test-bucket");
+ when(configuration.getPartSize()).thenReturn(25L * 1024 * 1024);
+ when(s3Client.putObject(any(PutObjectRequest.class),
any(RequestBody.class))).thenAnswer(invocation -> {
+ PutObjectRequest request = invocation.getArgument(0);
+ RequestBody body = invocation.getArgument(1);
+ uploadedLength = request.contentLength();
+ try (InputStream is = body.contentStreamProvider().newStream()) {
+ uploadedContent = is.readAllBytes();
+ }
+ return PutObjectResponse.builder()
+
.sdkHttpResponse(SdkHttpResponse.builder().statusCode(200).build())
+ .build();
+ });
+
+ producer = new AWS2S3Producer(endpoint);
+
+ payload = new byte[PAYLOAD_SIZE];
+ new Random(42).nextBytes(payload);
+ }
+
+ @Test
+ public void uploadPathBodyWithoutContentLength(@TempDir Path dir) throws
Exception {
+ Path file = Files.write(dir.resolve("big.bin"), payload);
+
+ Exchange exchange = new DefaultExchange(camelContext);
+ exchange.getIn().setHeader(AWS2S3Constants.KEY, "big.bin");
+ exchange.getIn().setBody(file);
+
+ producer.process(exchange);
+
+ assertEquals(PAYLOAD_SIZE, uploadedLength);
+ assertArrayEquals(payload, uploadedContent);
+ }
+
+ @Test
+ public void uploadReaderCacheBodyWithNonAsciiText() throws Exception {
+ // the length of a reader cache is a number of characters, not bytes
+ String text = "\u00e6\u00f8\u00e5".repeat(1000);
+ byte[] expected = text.getBytes(StandardCharsets.UTF_8);
+
+ Exchange exchange = new DefaultExchange(camelContext);
+ exchange.setProperty(Exchange.CHARSET_NAME,
StandardCharsets.UTF_8.name());
+ exchange.getIn().setHeader(AWS2S3Constants.KEY, "text.txt");
+ exchange.getIn().setBody(new ReaderCache(text));
+
+ producer.process(exchange);
+
+ assertEquals(expected.length, uploadedLength);
+ assertArrayEquals(expected, uploadedContent);
+ }
+
+ @Test
+ public void multiPartUploadStreamBodyWithUnknownLength() throws Exception {
+ when(configuration.isMultiPartUpload()).thenReturn(true);
+
+ // a stream that does not support mark/reset, such as a network stream
+ InputStream stream = new FilterInputStream(new
ByteArrayInputStream(payload)) {
+ @Override
+ public boolean markSupported() {
+ return false;
+ }
+ };
+
+ Exchange exchange = new DefaultExchange(camelContext);
+ exchange.getIn().setHeader(AWS2S3Constants.KEY, "big.bin");
+ exchange.getIn().setBody(stream);
+
+ producer.process(exchange);
+
+ assertEquals(PAYLOAD_SIZE, uploadedLength);
+ assertArrayEquals(payload, uploadedContent);
+ }
+}
diff --git
a/components/camel-azure/camel-azure-storage-blob/src/main/java/org/apache/camel/component/azure/storage/blob/BlobBlock.java
b/components/camel-azure/camel-azure-storage-blob/src/main/java/org/apache/camel/component/azure/storage/blob/BlobBlock.java
index c89a397c2be9..c2c72ee18eed 100644
---
a/components/camel-azure/camel-azure-storage-blob/src/main/java/org/apache/camel/component/azure/storage/blob/BlobBlock.java
+++
b/components/camel-azure/camel-azure-storage-blob/src/main/java/org/apache/camel/component/azure/storage/blob/BlobBlock.java
@@ -16,13 +16,14 @@
*/
package org.apache.camel.component.azure.storage.blob;
-import java.io.BufferedInputStream;
+import java.io.ByteArrayInputStream;
import java.io.IOException;
import java.io.InputStream;
import java.util.UUID;
import com.azure.core.util.Base64Util;
import com.azure.storage.blob.models.Block;
+import org.apache.camel.support.PayloadHelper;
public final class BlobBlock {
private final InputStream blockStream;
@@ -38,11 +39,13 @@ public final class BlobBlock {
}
public static BlobBlock createBlobBlock(final String blockId, final
InputStream inputStream) throws IOException {
- InputStream is = inputStream;
- if (!is.markSupported()) {
- is = new BufferedInputStream(is);
+ long length = PayloadHelper.getLength(inputStream);
+ if (length < 0) {
+ // the block must be read to determine its length
+ byte[] data = inputStream.readAllBytes();
+ return createBlobBlock(blockId, data.length, new
ByteArrayInputStream(data));
}
- return createBlobBlock(blockId, BlobUtils.getInputStreamLength(is),
is);
+ return createBlobBlock(blockId, length, inputStream);
}
public static BlobBlock createBlobBlock(final String blockId, final long
size, final InputStream inputStream) {
diff --git
a/components/camel-azure/camel-azure-storage-blob/src/main/java/org/apache/camel/component/azure/storage/blob/BlobStreamAndLength.java
b/components/camel-azure/camel-azure-storage-blob/src/main/java/org/apache/camel/component/azure/storage/blob/BlobStreamAndLength.java
index f650431da603..268a28d321cc 100644
---
a/components/camel-azure/camel-azure-storage-blob/src/main/java/org/apache/camel/component/azure/storage/blob/BlobStreamAndLength.java
+++
b/components/camel-azure/camel-azure-storage-blob/src/main/java/org/apache/camel/component/azure/storage/blob/BlobStreamAndLength.java
@@ -24,7 +24,9 @@ import java.io.IOException;
import java.io.InputStream;
import org.apache.camel.Exchange;
+import org.apache.camel.Message;
import org.apache.camel.WrappedFile;
+import org.apache.camel.support.PayloadHelper;
public final class BlobStreamAndLength {
@@ -37,48 +39,42 @@ public final class BlobStreamAndLength {
this.streamLength = streamLength;
}
- @SuppressWarnings("rawtypes")
public static BlobStreamAndLength
createBlobStreamAndLengthFromExchangeBody(final Exchange exchange) throws
IOException {
- Object body = exchange.getIn().getBody();
- Long blobSize =
exchange.getIn().getHeader(BlobConstants.BLOB_UPLOAD_SIZE, () -> null,
Long.class);
- exchange.getIn().removeHeader(BlobConstants.BLOB_UPLOAD_SIZE); //
remove to avoid issues for further uploads
+ final Message message = exchange.getIn();
+ Object body = message.getBody();
+ Long blobSize = message.getHeader(BlobConstants.BLOB_UPLOAD_SIZE, ()
-> null, Long.class);
+ message.removeHeader(BlobConstants.BLOB_UPLOAD_SIZE); // remove to
avoid issues for further uploads
- if (body instanceof WrappedFile wf) {
- // Get file length from WrappedFile before unwrapping (works for
remote files like SFTP)
- if (blobSize == null) {
- blobSize = wf.getFileLength();
- }
- body = wf.getFile();
+ if (body instanceof WrappedFile<?> wf && wf.getFile() instanceof File
file) {
+ body = file;
}
-
- if (body instanceof InputStream) {
- InputStream is = (InputStream) body;
- if (blobSize == null && !is.markSupported()) {
- is = new BufferedInputStream(is);
- }
- return new BlobStreamAndLength(is, blobSize != null ? blobSize :
BlobUtils.getInputStreamLength(is));
- }
- if (body instanceof File) {
- return new BlobStreamAndLength(new BufferedInputStream(new
FileInputStream((File) body)), ((File) body).length());
+ if (body instanceof File file) {
+ return new BlobStreamAndLength(new BufferedInputStream(new
FileInputStream(file)), file.length());
}
- if (body instanceof byte[]) {
- return new BlobStreamAndLength(new ByteArrayInputStream((byte[])
body), ((byte[]) body).length);
+ if (body instanceof byte[] bytes) {
+ return new BlobStreamAndLength(new ByteArrayInputStream(bytes),
bytes.length);
}
- // try as input stream
- final InputStream inputStream
- =
exchange.getContext().getTypeConverter().tryConvertTo(InputStream.class,
exchange, body);
+ // the length of a wrapped file (such as a remote file from SFTP) is
known without reading it
+ long length = blobSize != null ? blobSize :
PayloadHelper.getBodyLength(message);
- if (inputStream == null) {
- // fallback to string based
+ InputStream is = body instanceof InputStream inputStream
+ ? inputStream
+ :
exchange.getContext().getTypeConverter().tryConvertTo(InputStream.class,
exchange, body);
+ if (is == null) {
throw new IllegalArgumentException("Unsupported blob type:" +
body.getClass().getName());
}
- InputStream is = inputStream;
- if (blobSize == null && !is.markSupported()) {
- is = new BufferedInputStream(is);
+ if (length < 0) {
+ length = PayloadHelper.getLength(is);
+ }
+ if (length < 0) {
+ // copy the data to determine the length, which uses stream
caching so big payloads
+ // are spooled to disk when spooling is enabled
+ is = PayloadHelper.cacheStream(exchange, is);
+ length = PayloadHelper.getLength(is);
}
- return new BlobStreamAndLength(is, blobSize != null ? blobSize :
BlobUtils.getInputStreamLength(is));
+ return new BlobStreamAndLength(is, length);
}
public InputStream getInputStream() {
diff --git
a/components/camel-azure/camel-azure-storage-blob/src/main/java/org/apache/camel/component/azure/storage/blob/BlobUtils.java
b/components/camel-azure/camel-azure-storage-blob/src/main/java/org/apache/camel/component/azure/storage/blob/BlobUtils.java
index 7d104fd9f401..46f02fad3366 100644
---
a/components/camel-azure/camel-azure-storage-blob/src/main/java/org/apache/camel/component/azure/storage/blob/BlobUtils.java
+++
b/components/camel-azure/camel-azure-storage-blob/src/main/java/org/apache/camel/component/azure/storage/blob/BlobUtils.java
@@ -16,14 +16,12 @@
*/
package org.apache.camel.component.azure.storage.blob;
-import java.io.ByteArrayInputStream;
-import java.io.FileInputStream;
import java.io.IOException;
import java.io.InputStream;
import org.apache.camel.Exchange;
import org.apache.camel.Message;
-import org.apache.camel.StreamCache;
+import org.apache.camel.support.PayloadHelper;
import org.apache.camel.util.ObjectHelper;
public final class BlobUtils {
@@ -35,35 +33,15 @@ public final class BlobUtils {
return ObjectHelper.isEmpty(exchange) ? null : exchange.getIn();
}
+ /**
+ * Gets the length of the input stream, without reading it.
+ *
+ * @return the length, or <tt>-1</tt> if the length cannot be
determined without reading the stream
+ * @deprecated use {@link PayloadHelper#getLength(Object)}
+ */
+ @Deprecated(since = "4.23.0")
public static long getInputStreamLength(InputStream is) throws IOException
{
- if (is instanceof StreamCache streamCache) {
- long len = streamCache.length();
- if (len > 0) {
- return len;
- }
- } else if (is instanceof FileInputStream fis) {
- return fis.getChannel().size();
- }
-
- if (!is.markSupported()) {
- return -1;
- }
- if (is instanceof ByteArrayInputStream) {
- return is.available();
- }
- long size = 0;
- try {
- is.mark(1024);
- int i = is.available();
- while (i > 0) {
- long skip = is.skip(i);
- size += skip;
- i = is.available();
- }
- } finally {
- is.reset();
- }
- return size;
+ return PayloadHelper.getLength(is);
}
public static String getContainerName(final BlobConfiguration
configuration, final Exchange exchange) {
diff --git
a/components/camel-azure/camel-azure-storage-blob/src/test/java/org/apache/camel/component/azure/storage/blob/BlobStreamAndLengthTest.java
b/components/camel-azure/camel-azure-storage-blob/src/test/java/org/apache/camel/component/azure/storage/blob/BlobStreamAndLengthTest.java
new file mode 100644
index 000000000000..0ff986877162
--- /dev/null
+++
b/components/camel-azure/camel-azure-storage-blob/src/test/java/org/apache/camel/component/azure/storage/blob/BlobStreamAndLengthTest.java
@@ -0,0 +1,52 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.component.azure.storage.blob;
+
+import java.io.InputStream;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.Random;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.support.DefaultExchange;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+class BlobStreamAndLengthTest extends CamelTestSupport {
+
+ @Test
+ void testPathBodyWithoutUploadSize(@TempDir Path dir) throws Exception {
+ // bigger than the buffer of a BufferedInputStream, which a mark/reset
based length probe cannot handle
+ byte[] payload = new byte[1024 * 1024];
+ new Random(42).nextBytes(payload);
+ Path file = Files.write(dir.resolve("big.bin"), payload);
+
+ Exchange exchange = new DefaultExchange(context);
+ exchange.getIn().setBody(file);
+
+ BlobStreamAndLength blob =
BlobStreamAndLength.createBlobStreamAndLengthFromExchangeBody(exchange);
+
+ assertEquals(payload.length, blob.getStreamLength());
+ try (InputStream is = blob.getInputStream()) {
+ assertArrayEquals(payload, is.readAllBytes());
+ }
+ }
+}
diff --git
a/components/camel-azure/camel-azure-storage-datalake/src/main/java/org/apache/camel/component/azure/storage/datalake/DataLakeUtils.java
b/components/camel-azure/camel-azure-storage-datalake/src/main/java/org/apache/camel/component/azure/storage/datalake/DataLakeUtils.java
index e0e289fc570b..7b24aa71c914 100644
---
a/components/camel-azure/camel-azure-storage-datalake/src/main/java/org/apache/camel/component/azure/storage/datalake/DataLakeUtils.java
+++
b/components/camel-azure/camel-azure-storage-datalake/src/main/java/org/apache/camel/component/azure/storage/datalake/DataLakeUtils.java
@@ -21,8 +21,8 @@ import java.io.InputStream;
import org.apache.camel.Exchange;
import org.apache.camel.Message;
+import org.apache.camel.support.PayloadHelper;
import org.apache.camel.util.ObjectHelper;
-import org.apache.commons.io.IOUtils;
public final class DataLakeUtils {
private DataLakeUtils() {
@@ -35,14 +35,14 @@ public final class DataLakeUtils {
return exchange.getIn();
}
+ /**
+ * Gets the length of the input stream, without reading it.
+ *
+ * @return the length, or <tt>-1</tt> if the length cannot be
determined without reading the stream
+ * @deprecated use {@link PayloadHelper#getLength(Object)}
+ */
+ @Deprecated(since = "4.23.0")
public static Long getInputStreamLength(final InputStream inputStream)
throws IOException {
- if (!inputStream.markSupported()) {
- throw new IllegalArgumentException("Inputstream with mark reset
support required");
- }
-
- final long length = IOUtils.toByteArray(inputStream).length;
- inputStream.reset();
-
- return length;
+ return PayloadHelper.getLength(inputStream);
}
}
diff --git
a/components/camel-azure/camel-azure-storage-datalake/src/main/java/org/apache/camel/component/azure/storage/datalake/FileStreamAndLength.java
b/components/camel-azure/camel-azure-storage-datalake/src/main/java/org/apache/camel/component/azure/storage/datalake/FileStreamAndLength.java
index ce6c210db692..f09f25b3b587 100644
---
a/components/camel-azure/camel-azure-storage-datalake/src/main/java/org/apache/camel/component/azure/storage/datalake/FileStreamAndLength.java
+++
b/components/camel-azure/camel-azure-storage-datalake/src/main/java/org/apache/camel/component/azure/storage/datalake/FileStreamAndLength.java
@@ -24,7 +24,9 @@ import java.io.IOException;
import java.io.InputStream;
import org.apache.camel.Exchange;
+import org.apache.camel.Message;
import org.apache.camel.WrappedFile;
+import org.apache.camel.support.PayloadHelper;
public final class FileStreamAndLength {
private final InputStream inputStream;
@@ -35,42 +37,40 @@ public final class FileStreamAndLength {
this.streamLength = streamLength;
}
- @SuppressWarnings("rawtypes")
public static FileStreamAndLength
createFileStreamAndLengthFromExchangeBody(final Exchange exchange) throws
IOException {
- Object body = exchange.getIn().getBody();
- long fileLength = -1;
+ final Message message = exchange.getIn();
+ Object body = message.getBody();
- if (body instanceof WrappedFile wf) {
- // Get file length from WrappedFile before unwrapping (works for
remote files like SFTP)
- fileLength = wf.getFileLength();
- body = wf.getFile();
+ if (body instanceof WrappedFile<?> wf && wf.getFile() instanceof File
file) {
+ body = file;
}
-
- if (body instanceof InputStream) {
- if (!((InputStream) body).markSupported()) {
- throw new IllegalArgumentException("Inputstream does not
support mark rest operations");
- }
- // Use cached file length if available, otherwise calculate from
stream
- long length = fileLength > 0 ? fileLength :
DataLakeUtils.getInputStreamLength((InputStream) body);
- return new FileStreamAndLength((InputStream) body, length);
- }
-
- if (body instanceof File) {
- return new FileStreamAndLength(new BufferedInputStream(new
FileInputStream((File) body)), ((File) body).length());
+ if (body instanceof File file) {
+ return new FileStreamAndLength(new BufferedInputStream(new
FileInputStream(file)), file.length());
}
-
- if (body instanceof byte[]) {
- return new FileStreamAndLength(new ByteArrayInputStream((byte[])
body), ((byte[]) body).length);
+ if (body instanceof byte[] bytes) {
+ return new FileStreamAndLength(new ByteArrayInputStream(bytes),
bytes.length);
}
- final InputStream inputStream
- =
exchange.getContext().getTypeConverter().tryConvertTo(InputStream.class,
exchange, body);
+ // the length of a wrapped file (such as a remote file from SFTP) is
known without reading it
+ long length = PayloadHelper.getBodyLength(message);
- if (inputStream == null) {
+ InputStream is = body instanceof InputStream inputStream
+ ? inputStream
+ :
exchange.getContext().getTypeConverter().tryConvertTo(InputStream.class,
exchange, body);
+ if (is == null) {
throw new IllegalArgumentException("Unsupported file type");
}
- return new FileStreamAndLength(inputStream,
DataLakeUtils.getInputStreamLength(inputStream));
+ if (length < 0) {
+ length = PayloadHelper.getLength(is);
+ }
+ if (length < 0) {
+ // copy the data to determine the length, which uses stream
caching so big payloads
+ // are spooled to disk when spooling is enabled
+ is = PayloadHelper.cacheStream(exchange, is);
+ length = PayloadHelper.getLength(is);
+ }
+ return new FileStreamAndLength(is, length);
}
public InputStream getInputStream() {
diff --git
a/components/camel-azure/camel-azure-storage-datalake/src/test/java/org/apache/camel/component/azure/storage/datalake/FileStreamAndLengthTest.java
b/components/camel-azure/camel-azure-storage-datalake/src/test/java/org/apache/camel/component/azure/storage/datalake/FileStreamAndLengthTest.java
new file mode 100644
index 000000000000..633d668538b7
--- /dev/null
+++
b/components/camel-azure/camel-azure-storage-datalake/src/test/java/org/apache/camel/component/azure/storage/datalake/FileStreamAndLengthTest.java
@@ -0,0 +1,57 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.component.azure.storage.datalake;
+
+import java.io.ByteArrayInputStream;
+import java.io.FilterInputStream;
+import java.io.InputStream;
+import java.util.Random;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.support.DefaultExchange;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+class FileStreamAndLengthTest extends CamelTestSupport {
+
+ @Test
+ void testStreamBodyWithoutMarkSupport() throws Exception {
+ byte[] payload = new byte[64 * 1024];
+ new Random(42).nextBytes(payload);
+
+ // a stream that does not support mark/reset, such as a network stream
+ InputStream stream = new FilterInputStream(new
ByteArrayInputStream(payload)) {
+ @Override
+ public boolean markSupported() {
+ return false;
+ }
+ };
+
+ Exchange exchange = new DefaultExchange(context);
+ exchange.getIn().setBody(stream);
+
+ FileStreamAndLength file =
FileStreamAndLength.createFileStreamAndLengthFromExchangeBody(exchange);
+
+ assertEquals(payload.length, file.getStreamLength());
+ try (InputStream is = file.getInputStream()) {
+ assertArrayEquals(payload, is.readAllBytes());
+ }
+ }
+}
diff --git
a/components/camel-file/src/main/java/org/apache/camel/component/file/GenericFileHelper.java
b/components/camel-file/src/main/java/org/apache/camel/component/file/GenericFileHelper.java
index 278a93c88441..3e2c4c81b862 100644
---
a/components/camel-file/src/main/java/org/apache/camel/component/file/GenericFileHelper.java
+++
b/components/camel-file/src/main/java/org/apache/camel/component/file/GenericFileHelper.java
@@ -16,16 +16,23 @@
*/
package org.apache.camel.component.file;
+import java.io.ByteArrayInputStream;
import java.io.File;
import java.io.IOException;
+import java.io.InputStream;
+import java.io.Reader;
+import java.nio.charset.Charset;
import java.nio.file.Files;
import java.nio.file.LinkOption;
import java.nio.file.Path;
import java.util.function.Supplier;
import org.apache.camel.Exchange;
+import org.apache.camel.InvalidPayloadException;
+import org.apache.camel.Message;
import org.apache.camel.support.MessageHelper;
import org.apache.camel.util.FileUtil;
+import org.apache.camel.util.IOHelper;
public final class GenericFileHelper {
@@ -134,6 +141,30 @@ public final class GenericFileHelper {
return compactTarget.equals(dir) || compactTarget.startsWith(dir +
separator);
}
+ /**
+ * Gets the message body as an input stream with the content encoded in
the given charset.
+ * <p/>
+ * The body is read as a {@link Reader} (which decodes the content the
same way as converting the body to a
+ * {@link String}) and encoded while the stream is read, so a big body is
not loaded into memory. A {@link String}
+ * body, or a body that cannot be read as a {@link Reader}, is encoded as
a whole.
+ *
+ * @param exchange the exchange
+ * @param charset the charset to encode the content with
+ * @return the input stream
+ * @throws InvalidPayloadException if the body cannot be converted
+ * @throws IOException if the charset is not supported
+ */
+ public static InputStream toInputStream(Exchange exchange, String charset)
throws InvalidPayloadException, IOException {
+ Message message = exchange.getIn();
+ if (!(message.getBody() instanceof String)) {
+ Reader reader = message.getBody(Reader.class);
+ if (reader != null) {
+ return new IOHelper.ReaderInputStream(reader,
Charset.forName(charset));
+ }
+ }
+ return new
ByteArrayInputStream(message.getMandatoryBody(String.class).getBytes(charset));
+ }
+
public static String asExclusiveReadLockKey(GenericFile file, String key) {
// use the copy from absolute path as that was the original path of the
// file when the lock was acquired
diff --git
a/components/camel-file/src/test/java/org/apache/camel/component/file/GenericFileHelperTest.java
b/components/camel-file/src/test/java/org/apache/camel/component/file/GenericFileHelperTest.java
index 65bbb25d30ed..70ec1864acf7 100644
---
a/components/camel-file/src/test/java/org/apache/camel/component/file/GenericFileHelperTest.java
+++
b/components/camel-file/src/test/java/org/apache/camel/component/file/GenericFileHelperTest.java
@@ -16,17 +16,27 @@
*/
package org.apache.camel.component.file;
+import java.io.ByteArrayInputStream;
import java.io.File;
import java.io.IOException;
+import java.io.InputStream;
+import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.camel.support.DefaultExchange;
+import org.apache.camel.util.IOHelper;
import org.junit.jupiter.api.Assumptions;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -127,4 +137,39 @@ public class GenericFileHelperTest {
assertFalse(GenericFileHelper.isWithinDirectory("..", "", '/'));
assertFalse(GenericFileHelper.isWithinDirectory("../secret.txt", "",
'/'));
}
+
+ @Test
+ public void shouldEncodeStreamBodyWithCharsetWhileReading() throws
Exception {
+ // longer than the internal character buffer, with characters that are
encoded differently in the charsets
+ String text = "\u00e6\u00f8\u00e5 \u00a9 ".repeat(2000);
+
+ try (CamelContext context = new DefaultCamelContext()) {
+ Exchange exchange = new DefaultExchange(context);
+ exchange.setProperty(Exchange.CHARSET_NAME,
StandardCharsets.UTF_8.name());
+ exchange.getIn().setBody(new
ByteArrayInputStream(text.getBytes(StandardCharsets.UTF_8)));
+
+ try (InputStream is = GenericFileHelper.toInputStream(exchange,
StandardCharsets.ISO_8859_1.name())) {
+ // the body is encoded while it is read, instead of being
loaded into memory as a whole
+ assertInstanceOf(IOHelper.ReaderInputStream.class, is);
+ assertArrayEquals(text.getBytes(StandardCharsets.ISO_8859_1),
is.readAllBytes());
+ }
+ }
+ }
+
+ @Test
+ public void shouldEncodeStreamBodyWithCharsetThatWritesByteOrderMark()
throws Exception {
+ // longer than the internal character buffer, so the body is encoded
in several chunks
+ String text = "a\u00e6".repeat(4000);
+
+ try (CamelContext context = new DefaultCamelContext()) {
+ Exchange exchange = new DefaultExchange(context);
+ exchange.setProperty(Exchange.CHARSET_NAME,
StandardCharsets.UTF_8.name());
+ exchange.getIn().setBody(new
ByteArrayInputStream(text.getBytes(StandardCharsets.UTF_8)));
+
+ try (InputStream is = GenericFileHelper.toInputStream(exchange,
StandardCharsets.UTF_16.name())) {
+ // the byte order mark is written once, at the start
+ assertArrayEquals(text.getBytes(StandardCharsets.UTF_16),
is.readAllBytes());
+ }
+ }
+ }
}
diff --git
a/components/camel-ftp/src/main/java/org/apache/camel/component/file/remote/FtpOperations.java
b/components/camel-ftp/src/main/java/org/apache/camel/component/file/remote/FtpOperations.java
index 1dc597eb7944..56965c73e4da 100644
---
a/components/camel-ftp/src/main/java/org/apache/camel/component/file/remote/FtpOperations.java
+++
b/components/camel-ftp/src/main/java/org/apache/camel/component/file/remote/FtpOperations.java
@@ -783,7 +783,7 @@ public class FtpOperations implements
RemoteFileOperations<FTPFile> {
if (charset != null) {
// charset configured so we must convert to the desired
// charset so we can write with encoding
- is = new
ByteArrayInputStream(exchange.getIn().getMandatoryBody(String.class).getBytes(charset));
+ is = GenericFileHelper.toInputStream(exchange, charset);
log.trace("Using InputStream {} with charset {}.", is,
charset);
} else {
is = exchange.getIn().getMandatoryBody(InputStream.class);
diff --git
a/components/camel-ftp/src/main/java/org/apache/camel/component/file/remote/SftpOperations.java
b/components/camel-ftp/src/main/java/org/apache/camel/component/file/remote/SftpOperations.java
index e921423868fa..36eff0fbd983 100644
---
a/components/camel-ftp/src/main/java/org/apache/camel/component/file/remote/SftpOperations.java
+++
b/components/camel-ftp/src/main/java/org/apache/camel/component/file/remote/SftpOperations.java
@@ -23,7 +23,6 @@ import java.io.FileOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
-import java.io.UnsupportedEncodingException;
import java.net.InetAddress;
import java.net.Socket;
import java.nio.charset.Charset;
@@ -1226,7 +1225,7 @@ public class SftpOperations implements
RemoteFileOperations<SftpRemoteFile> {
if (charset != null) {
// charset configured so we must convert to the desired
// charset so we can write with encoding
- is = new
ByteArrayInputStream(exchange.getIn().getMandatoryBody(String.class).getBytes(charset));
+ is = GenericFileHelper.toInputStream(exchange, charset);
LOG.trace("Using InputStream {} with charset {}.", is,
charset);
} else {
is = exchange.getIn().getMandatoryBody(InputStream.class);
@@ -1264,7 +1263,7 @@ public class SftpOperations implements
RemoteFileOperations<SftpRemoteFile> {
} catch (SftpException e) {
createResultHeadersFromExchange(e, exchange);
throw new GenericFileOperationFailedException("Cannot store file:
" + name, e);
- } catch (UnsupportedEncodingException | InvalidPayloadException e) {
+ } catch (IOException | InvalidPayloadException e) {
throw new GenericFileOperationFailedException("Cannot store file:
" + name, e);
} finally {
IOHelper.close(is, "store: " + name, LOG);
diff --git
a/components/camel-ftp/src/test/java/org/apache/camel/component/file/remote/sftp/integration/SftpProducerWithCharsetIT.java
b/components/camel-ftp/src/test/java/org/apache/camel/component/file/remote/sftp/integration/SftpProducerWithCharsetIT.java
index a32a063702b6..4f7677abe3c0 100644
---
a/components/camel-ftp/src/test/java/org/apache/camel/component/file/remote/sftp/integration/SftpProducerWithCharsetIT.java
+++
b/components/camel-ftp/src/test/java/org/apache/camel/component/file/remote/sftp/integration/SftpProducerWithCharsetIT.java
@@ -16,7 +16,9 @@
*/
package org.apache.camel.component.file.remote.sftp.integration;
+import java.io.ByteArrayInputStream;
import java.io.File;
+import java.nio.charset.StandardCharsets;
import org.apache.camel.Exchange;
import org.apache.commons.io.FileUtils;
@@ -51,6 +53,22 @@ public class SftpProducerWithCharsetIT extends
SftpServerTestSupport {
assertEquals(SAMPLE_FILE_PAYLOAD, storedPayload);
}
+ @Test
+ public void testProducerWithCharsetFromStream() throws Exception {
+ // a stream body (such as from another component) is encoded to the
charset while it is uploaded
+ template.send(getSftpUri(), exchange -> {
+ exchange.setProperty(Exchange.CHARSET_NAME,
StandardCharsets.UTF_8.name());
+ exchange.getIn().setHeader(Exchange.FILE_NAME, SAMPLE_FILE_NAME);
+ exchange.getIn().setBody(new
ByteArrayInputStream(SAMPLE_FILE_PAYLOAD.getBytes(StandardCharsets.UTF_8)));
+ });
+
+ File file = new File(service.getFtpRootDir() + "/" + SAMPLE_FILE_NAME);
+ assertTrue(file.exists(), "The uploaded file should exist");
+
+ String storedPayload = FileUtils.readFileToString(file,
SAMPLE_FILE_CHARSET);
+ assertEquals(SAMPLE_FILE_PAYLOAD, storedPayload);
+ }
+
private String getSftpUri() {
return
"sftp://localhost:{{ftp.server.port}}/{{ftp.root.dir}}?username=admin&password=admin&charset="
+ SAMPLE_FILE_CHARSET + "&knownHostsFile=" +
service.getKnownHostsFile();
diff --git
a/components/camel-google/camel-google-storage/src/main/java/org/apache/camel/component/google/storage/GoogleCloudStorageProducer.java
b/components/camel-google/camel-google-storage/src/main/java/org/apache/camel/component/google/storage/GoogleCloudStorageProducer.java
index e876cf7cb56b..0de84d1c76ba 100644
---
a/components/camel-google/camel-google-storage/src/main/java/org/apache/camel/component/google/storage/GoogleCloudStorageProducer.java
+++
b/components/camel-google/camel-google-storage/src/main/java/org/apache/camel/component/google/storage/GoogleCloudStorageProducer.java
@@ -16,8 +16,6 @@
*/
package org.apache.camel.component.google.storage;
-import java.io.ByteArrayInputStream;
-import java.io.ByteArrayOutputStream;
import java.io.File;
import java.io.FileInputStream;
import java.io.IOException;
@@ -43,6 +41,7 @@ import org.apache.camel.Message;
import org.apache.camel.RuntimeCamelException;
import org.apache.camel.WrappedFile;
import org.apache.camel.support.DefaultProducer;
+import org.apache.camel.support.PayloadHelper;
import org.apache.camel.util.IOHelper;
import org.apache.camel.util.ObjectHelper;
import org.slf4j.Logger;
@@ -125,7 +124,7 @@ public class GoogleCloudStorageProducer extends
DefaultProducer {
objectMetadata.put("Content-Length", String.valueOf(fileLength));
}
// Handle Content-Length if not already set
- is = setContentLength(objectMetadata, is);
+ is = setContentLength(exchange, objectMetadata, obj, is);
Blob createdBlob;
BlobId blobId = BlobId.of(bucketName, objectName);
@@ -162,35 +161,35 @@ public class GoogleCloudStorageProducer extends
DefaultProducer {
}
/**
- * If no content-length header was found, calculate length by reading the
content.
+ * If no content-length header was found, determine the length of the
content.
*
+ * @param exchange the exchange
* @param objectMetadata Metadata set from Exchange headers
+ * @param body the Exchange body
* @param is InputStream to read the Exchange body content
- * @return the original InputStream if Content-Length is
set or a ByteArrayInputStream if the
- * original stream was read to determine the length.
+ * @return the original InputStream, or a copy of it if the
stream had to be copied to determine the
+ * length
* @throws IOException if the InputStream cannot be read.
*/
- private InputStream setContentLength(Map<String, String> objectMetadata,
InputStream is) throws IOException {
+ private InputStream setContentLength(Exchange exchange, Map<String,
String> objectMetadata, Object body, InputStream is)
+ throws IOException {
if (!objectMetadata.containsKey(Exchange.CONTENT_LENGTH) ||
objectMetadata.get(Exchange.CONTENT_LENGTH).equals("0")) {
- LOG.debug(
- "The content length is not defined. It needs to be
determined by reading the data into memory");
- ByteArrayOutputStream baos = determineLengthInputStream(is);
- objectMetadata.put("Content-Length", String.valueOf(baos.size()));
- return new ByteArrayInputStream(baos.toByteArray());
- } else {
- return is;
- }
- }
-
- private ByteArrayOutputStream determineLengthInputStream(InputStream is)
throws IOException {
- ByteArrayOutputStream out = new ByteArrayOutputStream();
- byte[] bytes = new byte[1024];
- int count;
- while ((count = is.read(bytes)) > 0) {
- out.write(bytes, 0, count);
+ // such as a java.nio.file.Path, byte[] or stream cache
+ long length = PayloadHelper.getLength(body);
+ if (length < 0) {
+ length = PayloadHelper.getLength(is);
+ }
+ if (length < 0) {
+ // copy the data to determine the length, which uses stream
caching
+ // so big payloads are spooled to disk when spooling is enabled
+ LOG.debug("The content length is not defined. It needs to be
determined by copying the data");
+ is = PayloadHelper.cacheStream(exchange, is);
+ length = PayloadHelper.getLength(is);
+ }
+ objectMetadata.put(Exchange.CONTENT_LENGTH,
String.valueOf(length));
}
- return out;
+ return is;
}
private Map<String, String> determineMetadata(final Exchange exchange) {
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 81d472d76d7a..94794597fc74 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
@@ -41,6 +41,7 @@ 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.support.PayloadHelper;
import org.apache.camel.util.FileUtil;
import org.apache.camel.util.ObjectHelper;
import org.slf4j.Logger;
@@ -144,16 +145,18 @@ public class IBMCOSProducer extends DefaultProducer {
} else if (wrappedFileLength > 0) {
// Use file length from WrappedFile (avoids reading stream into
memory)
metadata.setContentLength(wrappedFileLength);
- } else if (inputStream.markSupported()) {
- // For ByteArrayInputStream and similar streams that support
mark/reset
- inputStream.mark(Integer.MAX_VALUE);
- long length = 0;
- byte[] buffer = new byte[8192];
- int read;
- while ((read = inputStream.read(buffer)) != -1) {
- length += read;
+ } else {
+ // such as a java.nio.file.Path, byte[] or stream cache
+ long length = PayloadHelper.getLength(body);
+ if (length < 0) {
+ length = PayloadHelper.getLength(inputStream);
+ }
+ if (length < 0) {
+ // copy the data to calculate the length, which uses stream
caching
+ // so big payloads are spooled to disk when spooling is enabled
+ inputStream = PayloadHelper.cacheStream(exchange, inputStream);
+ length = PayloadHelper.getLength(inputStream);
}
- inputStream.reset();
metadata.setContentLength(length);
}
diff --git
a/components/camel-minio/src/main/java/org/apache/camel/component/minio/MinioProducer.java
b/components/camel-minio/src/main/java/org/apache/camel/component/minio/MinioProducer.java
index 456f614d8773..8e0e95ce54e4 100644
---
a/components/camel-minio/src/main/java/org/apache/camel/component/minio/MinioProducer.java
+++
b/components/camel-minio/src/main/java/org/apache/camel/component/minio/MinioProducer.java
@@ -16,8 +16,6 @@
*/
package org.apache.camel.component.minio;
-import java.io.ByteArrayInputStream;
-import java.io.ByteArrayOutputStream;
import java.io.File;
import java.io.FileInputStream;
import java.io.IOException;
@@ -48,6 +46,7 @@ import org.apache.camel.InvalidPayloadException;
import org.apache.camel.Message;
import org.apache.camel.WrappedFile;
import org.apache.camel.support.DefaultProducer;
+import org.apache.camel.support.PayloadHelper;
import org.apache.camel.util.FileUtil;
import org.apache.camel.util.IOHelper;
import org.slf4j.Logger;
@@ -157,18 +156,19 @@ public class MinioProducer extends DefaultProducer {
contentLength = filePayload.length();
}
} else {
+ if (contentLength <= 0) {
+ // such as a java.nio.file.Path, byte[] or stream cache
+ contentLength = PayloadHelper.getLength(object);
+ }
inputStream =
exchange.getMessage().getMandatoryBody(InputStream.class);
if (contentLength <= 0) {
- contentLength =
determineLengthInputStream(inputStream);
+ contentLength = PayloadHelper.getLength(inputStream);
if (contentLength == -1) {
- // fallback to read into memory to calculate length
- LOG.debug(
- "The content length is not defined. It
needs to be determined by reading the data into memory");
- ByteArrayOutputStream baos = new
ByteArrayOutputStream();
- IOHelper.copyAndCloseInput(inputStream, baos);
- byte[] arr = baos.toByteArray();
- contentLength = arr.length;
- inputStream = new ByteArrayInputStream(arr);
+ // fallback to copy the data to calculate the
length, which uses stream caching
+ // so big payloads are spooled to disk when
spooling is enabled
+ LOG.debug("The content length is not defined. It
needs to be determined by copying the data");
+ inputStream = PayloadHelper.cacheStream(exchange,
inputStream);
+ contentLength =
PayloadHelper.getLength(inputStream);
}
}
}
@@ -519,28 +519,6 @@ public class MinioProducer extends DefaultProducer {
return storageClass;
}
- private long determineLengthInputStream(InputStream is) throws IOException
{
- if (!is.markSupported()) {
- return -1;
- }
- if (is instanceof ByteArrayInputStream) {
- return is.available();
- }
- long size = 0;
- try {
- is.mark(MinioConstants.BYTE_ARRAY_LENGTH);
- int i = is.available();
- while (i > 0) {
- long skip = is.skip(i);
- size += skip;
- i = is.available();
- }
- } finally {
- is.reset();
- }
- return size;
- }
-
protected MinioConfiguration getConfiguration() {
return getEndpoint().getConfiguration();
}
diff --git
a/components/camel-smb/src/main/java/org/apache/camel/component/smb/SmbOperations.java
b/components/camel-smb/src/main/java/org/apache/camel/component/smb/SmbOperations.java
index 3ef511cf7b83..82827177cd09 100644
---
a/components/camel-smb/src/main/java/org/apache/camel/component/smb/SmbOperations.java
+++
b/components/camel-smb/src/main/java/org/apache/camel/component/smb/SmbOperations.java
@@ -493,7 +493,7 @@ public class SmbOperations implements SmbFileOperations {
if (charset != null) {
// charset configured so we must convert to the desired
// charset so we can write with encoding
- is = new
ByteArrayInputStream(exchange.getIn().getMandatoryBody(String.class).getBytes(charset));
+ is = GenericFileHelper.toInputStream(exchange, charset);
LOG.trace("Using InputStream {} with charset {}.", is,
charset);
} else {
is = exchange.getIn().getMandatoryBody(InputStream.class);
diff --git
a/core/camel-core/src/test/java/org/apache/camel/support/PayloadHelperTest.java
b/core/camel-core/src/test/java/org/apache/camel/support/PayloadHelperTest.java
new file mode 100644
index 000000000000..af727b45de40
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/support/PayloadHelperTest.java
@@ -0,0 +1,124 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.support;
+
+import java.io.BufferedInputStream;
+import java.io.ByteArrayInputStream;
+import java.io.FileInputStream;
+import java.io.InputStream;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.Random;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.StreamCache;
+import org.apache.camel.component.file.GenericFile;
+import org.apache.camel.converter.stream.FileInputStreamCache;
+import org.apache.camel.impl.engine.DefaultUnitOfWork;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+
+public class PayloadHelperTest extends ContextTestSupport {
+
+ private static final int PAYLOAD_SIZE = 256 * 1024;
+
+ private byte[] payload;
+ private Path file;
+
+ @Override
+ protected CamelContext createCamelContext() throws Exception {
+ CamelContext context = super.createCamelContext();
+ context.setStreamCaching(true);
+
context.getStreamCachingStrategy().setSpoolDirectory(testDirectory("spool").toFile());
+ context.getStreamCachingStrategy().setSpoolEnabled(true);
+ return context;
+ }
+
+ @Override
+ public boolean isUseRouteBuilder() {
+ return false;
+ }
+
+ @Override
+ @BeforeEach
+ public void setUp() throws Exception {
+ super.setUp();
+ payload = new byte[PAYLOAD_SIZE];
+ new Random(42).nextBytes(payload);
+ file = Files.write(testFile("payload.bin"), payload);
+ }
+
+ @Test
+ public void testLengthIsDeterminedWithoutReading() throws Exception {
+ assertEquals(PAYLOAD_SIZE, PayloadHelper.getLength(file));
+ assertEquals(PAYLOAD_SIZE, PayloadHelper.getLength(file.toFile()));
+ assertEquals(PAYLOAD_SIZE, PayloadHelper.getLength(payload));
+ assertEquals(PAYLOAD_SIZE, PayloadHelper.getLength(new
ByteArrayInputStream(payload)));
+ try (FileInputStream fis = new FileInputStream(file.toFile())) {
+ // only the remaining bytes are counted
+ fis.skip(1024);
+ assertEquals(PAYLOAD_SIZE - 1024, PayloadHelper.getLength(fis));
+ }
+
+ // the length of a buffered stream is unknown, and determining it must
not read the stream
+ try (InputStream is = new
BufferedInputStream(Files.newInputStream(file))) {
+ assertEquals(-1, PayloadHelper.getLength(is));
+ assertArrayEquals(payload, is.readAllBytes());
+ }
+ }
+
+ @Test
+ public void testBodyLengthOfWrappedFile() {
+ Exchange exchange =
context.getEndpoint("mock:result").createExchange();
+
+ // a remote file whose content is not stored in a local file, such as
from SFTP
+ GenericFile<Object> remote = new GenericFile<>();
+ remote.setFile(new Object());
+ remote.setFileLength(PAYLOAD_SIZE);
+ exchange.getIn().setBody(remote);
+ assertEquals(PAYLOAD_SIZE,
PayloadHelper.getBodyLength(exchange.getIn()));
+
+ // a local file without a known file length
+ GenericFile<Object> local = new GenericFile<>();
+ local.setFile(file.toFile());
+ exchange.getIn().setBody(local);
+ assertEquals(PAYLOAD_SIZE,
PayloadHelper.getBodyLength(exchange.getIn()));
+ }
+
+ @Test
+ public void testCacheStreamSpoolsBigPayloadToDisk() throws Exception {
+ context.start();
+ Exchange exchange =
context.getEndpoint("mock:result").createExchange();
+ exchange.getExchangeExtension().setUnitOfWork(new
DefaultUnitOfWork(exchange));
+
+ // a stream of unknown length
+ InputStream stream = new BufferedInputStream(new
ByteArrayInputStream(payload));
+ InputStream cached = PayloadHelper.cacheStream(exchange, stream);
+
+ // the payload is above the spool threshold, so it is kept on disk and
not in memory
+ assertInstanceOf(FileInputStreamCache.class, cached);
+ assertEquals(PAYLOAD_SIZE, ((StreamCache) cached).length());
+ assertEquals(PAYLOAD_SIZE, PayloadHelper.getLength(cached));
+ assertArrayEquals(payload, cached.readAllBytes());
+ }
+}
diff --git
a/core/camel-support/src/main/java/org/apache/camel/support/PayloadHelper.java
b/core/camel-support/src/main/java/org/apache/camel/support/PayloadHelper.java
new file mode 100644
index 000000000000..187412e94630
--- /dev/null
+++
b/core/camel-support/src/main/java/org/apache/camel/support/PayloadHelper.java
@@ -0,0 +1,119 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.support;
+
+import java.io.ByteArrayInputStream;
+import java.io.File;
+import java.io.FileInputStream;
+import java.io.IOException;
+import java.io.InputStream;
+import java.nio.channels.FileChannel;
+import java.nio.file.Files;
+import java.nio.file.Path;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.Message;
+import org.apache.camel.StreamCache;
+import org.apache.camel.WrappedFile;
+import org.apache.camel.support.builder.OutputStreamBuilder;
+import org.apache.camel.util.IOHelper;
+
+/**
+ * Helper for components that need to know the length of a (possibly big)
message body, for example to upload it.
+ * <p/>
+ * The length is determined without reading the body. When the length cannot
be determined that way, the body can be
+ * copied with {@link #cacheStream(Exchange, InputStream)}, which uses stream
caching and therefore spools big payloads
+ * to disk when spooling is enabled, instead of loading them into memory.
+ */
+public final class PayloadHelper {
+
+ private PayloadHelper() {
+ }
+
+ /**
+ * Gets the length of the message body, without reading the body.
+ *
+ * @param message the message
+ * @return the length in bytes, or <tt>-1</tt> if the length
cannot be determined without reading the body
+ */
+ public static long getBodyLength(Message message) {
+ Object body = message.getBody();
+ if (body instanceof WrappedFile<?> wf) {
+ long length = wf.getFileLength();
+ if (length > 0) {
+ return length;
+ }
+ Object file = wf.getFile();
+ body = file instanceof File || file instanceof Path ? file :
wf.getBody();
+ }
+ return getLength(body);
+ }
+
+ /**
+ * Gets the length of the given value, without reading it.
+ * <p/>
+ * The length is known for files, byte arrays, stream caches of bytes, and
input streams backed by a byte array or a
+ * file.
+ *
+ * @param value the value, such as a message body or an input stream
+ * @return the length in bytes, or <tt>-1</tt> if the length cannot
be determined without reading the value
+ */
+ public static long getLength(Object value) {
+ try {
+ if (value instanceof File file) {
+ return file.isFile() ? file.length() : -1;
+ } else if (value instanceof Path path) {
+ return Files.isRegularFile(path) ? Files.size(path) : -1;
+ } else if (value instanceof byte[] bytes) {
+ return bytes.length;
+ } else if (value instanceof StreamCache cache && value instanceof
InputStream && cache.length() > 0) {
+ // only a stream cache that is a byte stream has a length in
bytes (a reader cache counts characters),
+ // and a stream cache may return 0 when the length cannot be
computed
+ return cache.length();
+ } else if (value instanceof ByteArrayInputStream is) {
+ return is.available();
+ } else if (value instanceof FileInputStream fis) {
+ FileChannel channel = fis.getChannel();
+ return channel.size() - channel.position();
+ }
+ } catch (IOException | UnsupportedOperationException e) {
+ // the length cannot be determined
+ }
+ return -1;
+ }
+
+ /**
+ * Copies the given input stream so its length is known.
+ * <p/>
+ * The stream is copied using stream caching (when enabled), so a big
payload is spooled to disk if spooling is
+ * enabled, instead of being loaded into memory. The copy is released when
the exchange is complete.
+ *
+ * @param exchange the exchange
+ * @param is the input stream to copy, which is closed afterwards
+ * @return the copy, which is a stream whose length can be
determined with {@link #getLength(Object)}
+ * @throws IOException if the stream cannot be copied
+ */
+ public static InputStream cacheStream(Exchange exchange, InputStream is)
throws IOException {
+ OutputStreamBuilder builder =
OutputStreamBuilder.withExchange(exchange);
+ IOHelper.copyAndCloseInput(is, builder);
+ Object answer = builder.build();
+ if (answer instanceof InputStream cached) {
+ return cached;
+ }
+ return new ByteArrayInputStream((byte[]) answer);
+ }
+}
diff --git a/core/camel-util/src/main/java/org/apache/camel/util/IOHelper.java
b/core/camel-util/src/main/java/org/apache/camel/util/IOHelper.java
index bb9d3ded681c..58ea4874c69d 100644
--- a/core/camel-util/src/main/java/org/apache/camel/util/IOHelper.java
+++ b/core/camel-util/src/main/java/org/apache/camel/util/IOHelper.java
@@ -39,10 +39,14 @@ import java.nio.channels.FileChannel;
import java.nio.channels.ReadableByteChannel;
import java.nio.channels.WritableByteChannel;
import java.nio.charset.Charset;
+import java.nio.charset.CharsetEncoder;
+import java.nio.charset.CoderResult;
+import java.nio.charset.CodingErrorAction;
import java.nio.charset.UnsupportedCharsetException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Locale;
+import java.util.Objects;
import java.util.Scanner;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
@@ -729,49 +733,87 @@ public final class IOHelper {
}
/**
- * Encoding-aware input stream.
+ * An input stream that reads characters from a {@link Reader} and encodes
them with the given charset, a chunk at a
+ * time, so the content is never held in memory as a whole.
*/
- public static class EncodingInputStream extends InputStream {
+ public static class ReaderInputStream extends InputStream {
private final Lock lock = new ReentrantLock();
- private final Path file;
- private final BufferedReader reader;
- private final Charset defaultStreamCharset;
+ private final Reader reader;
+ // a single encoder for the whole stream, so a charset that writes a
byte order mark (such as UTF-16) writes
+ // it only once, and a stateful charset keeps its state from one chunk
to the next
+ private final CharsetEncoder encoder;
+ private final CharBuffer chars = CharBuffer.allocate(4096);
+ private final ByteBuffer bytes;
+ private boolean endOfInput;
+ private boolean encoded;
+ private boolean flushed;
- private ByteBuffer bufferBytes;
- private final CharBuffer bufferedChars = CharBuffer.allocate(4096);
- // the first half of a surrogate pair that was read at the end of the
buffer
- private char pendingHighSurrogate;
-
- public EncodingInputStream(Path file, String charset) throws
IOException {
- this.file = file;
- reader = toReader(file, charset);
- defaultStreamCharset = defaultCharset.get();
+ /**
+ * @param reader the reader to read the characters from
+ * @param charset the charset to encode the characters with
+ */
+ public ReaderInputStream(Reader reader, Charset charset) {
+ this.reader = reader;
+ // replace malformed and unmappable characters, the same way as
Charset.encode and String.getBytes
+ this.encoder = charset.newEncoder()
+ .onMalformedInput(CodingErrorAction.REPLACE)
+ .onUnmappableCharacter(CodingErrorAction.REPLACE);
+ this.bytes = ByteBuffer.allocate((int) Math.ceil(chars.capacity()
* encoder.maxBytesPerChar()));
+ // nothing to read yet
+ bytes.limit(0);
}
@Override
public int read() throws IOException {
- while (bufferBytes == null || bufferBytes.remaining() <= 0) {
- BufferCaster.cast(bufferedChars).clear();
- if (pendingHighSurrogate != 0) {
- bufferedChars.put(pendingHighSurrogate);
- pendingHighSurrogate = 0;
- }
- int len = reader.read(bufferedChars);
- bufferedChars.flip();
- if (len == -1 && !bufferedChars.hasRemaining()) {
- return -1;
+ return fill() ? bytes.get() & 0xFF : -1;
+ }
+
+ @Override
+ public int read(byte[] b, int off, int len) throws IOException {
+ Objects.checkFromIndexSize(off, len, b.length);
+ if (len == 0) {
+ return 0;
+ }
+ if (!fill()) {
+ return -1;
+ }
+ int n = Math.min(len, bytes.remaining());
+ bytes.get(b, off, n);
+ return n;
+ }
+
+ /**
+ * Encodes the next chunk of characters if all the encoded bytes have
been read.
+ *
+ * @return <tt>false</tt> if the end of the reader has been reached
+ */
+ private boolean fill() throws IOException {
+ while (!bytes.hasRemaining()) {
+ if (flushed) {
+ return false;
}
- int limit = bufferedChars.limit();
- if (len != -1 && limit > 0 &&
Character.isHighSurrogate(bufferedChars.get(limit - 1))) {
- // a surrogate pair (such as an emoji) is split at the end
of the buffer, so encode the high
- // surrogate together with the low surrogate in the next
read (alone it would be encoded as ?)
- pendingHighSurrogate = bufferedChars.get(limit - 1);
- bufferedChars.limit(limit - 1);
+ bytes.clear();
+ if (encoded) {
+ // write what the encoder may still hold
+ flushed = encoder.flush(bytes).isUnderflow();
+ } else {
+ if (!endOfInput && reader.read(chars) == -1) {
+ endOfInput = true;
+ }
+ chars.flip();
+ // characters that cannot be encoded yet (such as the
first half of a surrogate pair at the end
+ // of the buffer) are kept for the next chunk
+ CoderResult result = encoder.encode(chars, bytes,
endOfInput);
+ chars.compact();
+ if (endOfInput && result.isUnderflow()) {
+ encoded = true;
+ flushed = encoder.flush(bytes).isUnderflow();
+ }
}
- bufferBytes = defaultStreamCharset.encode(bufferedChars);
+ bytes.flip();
}
- return bufferBytes.get() & 0xFF;
+ return true;
}
@Override
@@ -784,10 +826,30 @@ public final class IOHelper {
lock.lock();
try {
reader.reset();
+ encoder.reset();
+ chars.clear();
+ bytes.clear();
+ bytes.limit(0);
+ endOfInput = false;
+ encoded = false;
+ flushed = false;
} finally {
lock.unlock();
}
}
+ }
+
+ /**
+ * Encoding-aware input stream.
+ */
+ public static class EncodingInputStream extends ReaderInputStream {
+
+ private final Path file;
+
+ public EncodingInputStream(Path file, String charset) throws
IOException {
+ super(toReader(file, charset), defaultCharset.get());
+ this.file = file;
+ }
public InputStream toOriginalInputStream() throws IOException {
return Files.newInputStream(file);
diff --git a/docs/user-manual/modules/ROOT/nav.adoc
b/docs/user-manual/modules/ROOT/nav.adoc
index 82397594fc67..616903789582 100644
--- a/docs/user-manual/modules/ROOT/nav.adoc
+++ b/docs/user-manual/modules/ROOT/nav.adoc
@@ -105,6 +105,7 @@
** xref:routes.adoc[Routes]
** xref:startup-condition.adoc[Startup Condition]
** xref:stream-caching.adoc[Stream caching]
+*** xref:large-payloads.adoc[Large payloads]
** xref:template-engines.adoc[Template Engines]
** xref:transformer.adoc[Transformer]
** xref:threading-model.adoc[Threading Model]
diff --git a/docs/user-manual/modules/ROOT/pages/index.adoc
b/docs/user-manual/modules/ROOT/pages/index.adoc
index e05b5bac6721..7d325b6ef5bf 100644
--- a/docs/user-manual/modules/ROOT/pages/index.adoc
+++ b/docs/user-manual/modules/ROOT/pages/index.adoc
@@ -107,6 +107,7 @@ For a deeper and better understanding of Apache Camel, an
xref:faq:what-is-camel
* xref:route-diagram.adoc[Visual Route Diagrams]
* xref:routes.adoc[Routes]
* xref:stream-caching.adoc[Stream caching]
+** xref:large-payloads.adoc[Large payloads]
* xref:threading-model.adoc[Threading Model]
* xref:tracer.adoc[Tracer]
* xref:type-converter.adoc[Type Converter]
diff --git a/docs/user-manual/modules/ROOT/pages/large-payloads.adoc
b/docs/user-manual/modules/ROOT/pages/large-payloads.adoc
new file mode 100644
index 000000000000..a1eed080b73f
--- /dev/null
+++ b/docs/user-manual/modules/ROOT/pages/large-payloads.adoc
@@ -0,0 +1,205 @@
+= Large payloads
+
+Camel can move large files, such as several gigabytes, between components
without loading them into memory.
+This does not happen with every configuration, though: a few defaults and a
few route constructs load the whole
+message body into the heap. This page explains which ones, and how to
configure the most common components to
+stream large payloads.
+
+== Stream caching
+
+xref:stream-caching.adoc[Stream caching] is enabled by default, but spooling
to disk is *not*. When a message body
+is a stream, such as an `InputStream` returned by a component, stream caching
copies it into a re-readable cache
+the next time a processor needs the body. Without spooling, that cache is kept
in memory, so a stream of several
+gigabytes ends up in the heap.
+
+For routes that handle large payloads, choose one of the following:
+
+* Enable spooling, so bodies above the spool threshold (128 KB by default) are
cached in a temporary file:
++
+[source,properties]
+----
+camel.main.streamCachingSpoolEnabled = true
+# a volume with room for the largest payload times the number of concurrent
exchanges
+camel.main.streamCachingSpoolDirectory = /data/camel-spool
+----
++
+The body can then be read many times, for example by redelivery, and the heap
stays bounded. The cost is one full
+copy of the payload to the spool directory. In containers, make sure the spool
directory is writable and large
+enough.
+
+* Disable stream caching for the route, so the stream is passed from component
to component as-is:
++
+[source,java]
+----
+from("sftp:...?streamDownload=true")
+ .streamCache("false")
+ .to("http:...");
+----
++
+Nothing is copied, but the body can only be read once. Do not use steps that
read the body again, such as
+redelivery, a `choice` on the body, `multicast`, `wireTap` or `recipientList`.
+
+A body that is a file (`java.io.File`, `java.nio.file.Path`, or the file of
the xref:components::file-component.adoc[File]
+component) is not cached, as it can already be read many times. Keeping a
large payload as a file is often the
+cheapest option.
+
+== Route steps that load the whole body
+
+Even when the components stream, the following steps load the whole body into
memory:
+
+* converting the body to a `String` or `byte[]`, for example with
`convertBodyTo(String.class)`
+* using `+${body}+` in a
xref:components:languages:simple-language.adoc[Simple] expression, including
`+.log("${body}")+`
+* unmarshalling the whole body with a data format, such as JSON or XML
+* tracing the body with the xref:tracer.adoc[Tracer] or the
xref:backlog-tracer.adoc[Backlog Tracer]; the body is
+converted before it is clipped to the maximum number of characters
+* the xref:components:eips:split-eip.adoc[Split] EIP without `streaming()`
+
+To process a large file record by record, split it in streaming mode, for
example
+`split(body().tokenize("\n")).streaming()`.
+
+== Components
+
+The following table describes how the most common components handle large
payloads, and which options make them
+stream. When a stream is received, stream caching applies to it as described
above.
+
+[width="100%",cols="2,4,4",options="header"]
+|===
+| Component | Receiving (consumer, download) | Sending (producer, upload)
+
+| xref:components::file-component.adoc[File]
+| The body is the file; it is not read until it is used.
+| A `File` or `InputStream` body is streamed to the file. With `charset`, the
body is converted while it is written.
+
+| xref:components::ftp-component.adoc[FTP],
+xref:components::sftp-component.adoc[SFTP],
+xref:components::smb-component.adoc[SMB]
+| By default the whole file is loaded into memory. Set `streamDownload=true`
to receive an `InputStream`, or
+`localWorkDirectory` to download to a local file first (the body is then a
file).
+| A `File` or `InputStream` body is streamed. With `charset`, the body is
converted while it is uploaded.
+
+| xref:components::scp-component.adoc[SCP]
+| Not supported.
+| The whole body is loaded into memory.
+
+| xref:components::aws2-s3-component.adoc[AWS S3]
+| By default (`includeBody=true`) the whole object is loaded into memory. Set
`includeBody=false` to receive the
+object stream, and `autocloseBody=true` to close it when the exchange is done.
+| A file, or a stream whose length is known (a stream cache, or the
`CamelAwsS3ContentLength` header), is uploaded
+as a stream. Otherwise the body is copied first to find its length: with
stream caching and spooling enabled the
+copy is a temporary file, otherwise it is in memory. Use
`multiPartUpload=true` to upload large objects in parts.
+
+| xref:components::azure-storage-blob-component.adoc[Azure Storage Blob],
+xref:components::azure-storage-datalake-component.adoc[Azure Storage Data Lake]
+| The blob is received as a stream. With `fileDir`, the blob is downloaded to
a file.
+| Same as AWS S3. For blobs, the length can also be given with the
`CamelAzureStorageBlobUploadSize` header.
+
+| xref:components::google-storage-component.adoc[Google Storage]
+| With `includeBody=true` the object is copied into a stream cache (a
temporary file when spooling is enabled). With
+`downloadFileName`, the object is downloaded to a file.
+| Same as AWS S3. The length can also be given with the
`CamelGoogleCloudStorageContentLength` header.
+
+| xref:components::minio-component.adoc[MinIO],
+xref:components::ibm-cos-component.adoc[IBM Cloud Object Storage]
+| MinIO loads the whole object into memory when the body is included.
+| Same as AWS S3.
+
+| xref:components::http-component.adoc[HTTP]
+| The response is copied into a stream cache (a temporary file when spooling
is enabled). With
+`disableStreamCache=true` the producer returns the raw response stream, which
the route's stream caching then
+handles; disable stream caching on the route to pass it on without any copy.
+| A `File` or `InputStream` body is streamed. A `Content-Encoding: gzip`
request header compresses the whole body in
+memory.
+
+| xref:components::platform-http-component.adoc[Platform HTTP]
+| By default the whole request body is loaded into memory. Multipart file
uploads are written to temporary files,
+and a single uploaded file becomes the body. With `useStreaming=true` the
request body is written into a stream
+cache as it arrives (a temporary file when spooling is enabled), and the route
starts once the whole request is
+received.
+| A `File` or `InputStream` response body is streamed to the client.
+
+| xref:components::vertx-http-component.adoc[Vert.x HTTP]
+| The whole response is loaded into memory.
+| The whole request body is loaded into memory, unless it is a Vert.x
`ReadStream`.
+
+| xref:components::netty-http-component.adoc[Netty HTTP]
+| The whole response is loaded into memory, up to `chunkedMaxContentLength` (1
MB by default). With
+`disableStreamCache=true`, chunked responses are streamed.
+| With `disableStreamCache=true`, an `InputStream` body is sent chunked.
Otherwise it is loaded into memory.
+|===
+
+=== Platform HTTP request size limits
+
+The size of an uploaded request is limited by the runtime:
+
+* Camel Main and Camel JBang: `camel.server.maxBodySize`. When it is not set,
the Vert.x default of 10 MB applies.
+With `useStreaming=true` this limit does not apply, so limit the size of
requests in front of Camel if needed.
+* Quarkus: `quarkus.http.limits.max-body-size`, and the `quarkus.http.body.*`
options for file uploads. See the
+Quarkus documentation.
+* Spring Boot: `spring.servlet.multipart.max-file-size` and
`spring.servlet.multipart.max-request-size` for
+multipart uploads. See the Spring Boot documentation.
+
+On Camel Main, Camel JBang and Quarkus, `useStreaming=true` does not accept
multipart requests; use a separate
+endpoint for multipart uploads.
+
+== Examples
+
+=== From SFTP to HTTP
+
+Receive the remote file as a stream and send it on without copying it:
+
+[source,java]
+----
+from("sftp:host/inbox?username=...&streamDownload=true")
+ .streamCache("false")
+ .to("http:backend/upload");
+----
+
+The HTTP producer sends the stream with chunked transfer encoding.
Alternatively, keep stream caching and set
+`localWorkDirectory` on the SFTP endpoint, so the file is downloaded to disk
and the body is a file.
+
+=== From Platform HTTP to AWS S3
+
+Multipart uploads are written to temporary files by the HTTP server, and the
uploaded file is then sent to S3 from
+disk:
+
+[source,java]
+----
+from("platform-http:/upload?httpMethodRestrict=POST")
+ .setHeader(AWS2S3Constants.KEY, header(Exchange.FILE_NAME))
+ .to("aws2-s3:my-bucket");
+----
+
+Raise the request size limit of the runtime, as described above.
+
+For a raw request body (not multipart), enable spooling and use
`useStreaming=true`:
+
+[source,java]
+----
+from("platform-http:/upload?httpMethodRestrict=PUT&useStreaming=true")
+ .to("aws2-s3:my-bucket?keyName=upload.bin");
+----
+
+The request is written to the spool directory while it is received, and S3
uploads it from there with its known
+length.
+
+=== From AWS S3 to SFTP
+
+Receive the object as a stream and upload it without copying it:
+
+[source,java]
+----
+from("aws2-s3:my-bucket?includeBody=false&autocloseBody=true")
+ .streamCache("false")
+ .setHeader(Exchange.FILE_NAME, header(AWS2S3Constants.KEY))
+ .to("sftp:host/outbox?username=...");
+----
+
+=== Downloading a large file over HTTP
+
+[source,java]
+----
+from("timer:download?repeatCount=1")
+ .streamCache("false")
+ .to("http:server/big-file?disableStreamCache=true")
+ .to("file:downloads?fileName=big-file");
+----
diff --git a/docs/user-manual/modules/ROOT/pages/stream-caching.adoc
b/docs/user-manual/modules/ROOT/pages/stream-caching.adoc
index 61bf1a30ec84..4b8bdbf8e43b 100644
--- a/docs/user-manual/modules/ROOT/pages/stream-caching.adoc
+++ b/docs/user-manual/modules/ROOT/pages/stream-caching.adoc
@@ -6,6 +6,8 @@ Streams are cached in memory. However, for large stream
messages, you can set `s
and then large message (over 128 KB) will be cached in a temporary file
instead.
Camel itself will handle deleting the temporary file once the cached stream is
no longer necessary.
+TIP: For routes that move large payloads, such as big files, see
xref:large-payloads.adoc[Large payloads].
+
== Why is my message empty?
In Camel the message body can be of any types. Some types are safely