diff --git a/nifi-commons/nifi-site-to-site-client/src/main/java/org/apache/nifi/remote/codec/StandardFlowFileCodec.java b/nifi-commons/nifi-site-to-site-client/src/main/java/org/apache/nifi/remote/codec/StandardFlowFileCodec.java
index 1e998be33dcf..7d9caefea1fb 100644
--- a/nifi-commons/nifi-site-to-site-client/src/main/java/org/apache/nifi/remote/codec/StandardFlowFileCodec.java
+++ b/nifi-commons/nifi-site-to-site-client/src/main/java/org/apache/nifi/remote/codec/StandardFlowFileCodec.java
@@ -60,7 +60,7 @@ public void encode(final DataPacket dataPacket, final OutputStream encodedOut) t
out.writeLong(dataPacket.getSize());
final InputStream in = dataPacket.getData();
- StreamUtils.copy(in, encodedOut);
+ in.transferTo(encodedOut);
}
@Override
diff --git a/nifi-commons/nifi-site-to-site-client/src/main/java/org/apache/nifi/remote/util/SiteToSiteRestApiClient.java b/nifi-commons/nifi-site-to-site-client/src/main/java/org/apache/nifi/remote/util/SiteToSiteRestApiClient.java
index 5ea962d2410f..2efb95d1e7f2 100644
--- a/nifi-commons/nifi-site-to-site-client/src/main/java/org/apache/nifi/remote/util/SiteToSiteRestApiClient.java
+++ b/nifi-commons/nifi-site-to-site-client/src/main/java/org/apache/nifi/remote/util/SiteToSiteRestApiClient.java
@@ -37,7 +37,6 @@
import org.apache.nifi.remote.protocol.http.HttpHeaders;
import org.apache.nifi.remote.protocol.http.HttpProxy;
import org.apache.nifi.security.cert.StandardPrincipalFormatter;
-import org.apache.nifi.stream.io.StreamUtils;
import org.apache.nifi.web.api.dto.ControllerDTO;
import org.apache.nifi.web.api.dto.remote.PeerDTO;
import org.apache.nifi.web.api.entity.ControllerEntity;
@@ -505,7 +504,7 @@ private IOException handleErrResponse(final int responseCode, final InputStream
private TransactionResultEntity readResponse(final InputStream inputStream) throws IOException {
final ByteArrayOutputStream bos = new ByteArrayOutputStream();
- StreamUtils.copy(inputStream, bos);
+ inputStream.transferTo(bos);
String responseMessage = null;
try {
diff --git a/nifi-commons/nifi-site-to-site-client/src/test/java/org/apache/nifi/remote/protocol/SiteToSiteTestUtils.java b/nifi-commons/nifi-site-to-site-client/src/test/java/org/apache/nifi/remote/protocol/SiteToSiteTestUtils.java
index aea1485e5d2e..1a6849befe9e 100644
--- a/nifi-commons/nifi-site-to-site-client/src/test/java/org/apache/nifi/remote/protocol/SiteToSiteTestUtils.java
+++ b/nifi-commons/nifi-site-to-site-client/src/test/java/org/apache/nifi/remote/protocol/SiteToSiteTestUtils.java
@@ -19,7 +19,6 @@
import org.apache.nifi.remote.Transaction;
import org.apache.nifi.remote.TransactionCompletion;
import org.apache.nifi.remote.util.StandardDataPacket;
-import org.apache.nifi.stream.io.StreamUtils;
import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
@@ -43,7 +42,7 @@ public static DataPacket createDataPacket(String contents) {
public static String readContents(DataPacket packet) throws IOException {
ByteArrayOutputStream os = new ByteArrayOutputStream((int) packet.getSize());
- StreamUtils.copy(packet.getData(), os);
+ packet.getData().transferTo(os);
return os.toString(StandardCharsets.UTF_8);
}
diff --git a/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/stream/io/StreamUtils.java b/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/stream/io/StreamUtils.java
index 7d97045a5a98..89ab54d554a7 100644
--- a/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/stream/io/StreamUtils.java
+++ b/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/stream/io/StreamUtils.java
@@ -28,6 +28,15 @@
public class StreamUtils {
+ /**
+ * Copies from source to destination.
+ * @param source source InputStream
+ * @param destination destination OutputStream
+ * @return Total number of bytes copied
+ * @throws IOException If an error occurs when copying.
+ * @deprecated Use {@link InputStream#transferTo(OutputStream)} instead.
+ */
+ @Deprecated(since = "2.12.0", forRemoval = true)
public static long copy(final InputStream source, final OutputStream destination) throws IOException {
final byte[] buffer = new byte[8192];
int len;
diff --git a/nifi-commons/nifi-utils/src/test/java/org/apache/nifi/stream/io/LimitingInputStreamTest.java b/nifi-commons/nifi-utils/src/test/java/org/apache/nifi/stream/io/LimitingInputStreamTest.java
index 5dcefe4a28e1..cc69791a1a5f 100644
--- a/nifi-commons/nifi-utils/src/test/java/org/apache/nifi/stream/io/LimitingInputStreamTest.java
+++ b/nifi-commons/nifi-utils/src/test/java/org/apache/nifi/stream/io/LimitingInputStreamTest.java
@@ -33,7 +33,7 @@ public class LimitingInputStreamTest {
@Test
public void testReadLimitNotReached() throws IOException {
final LimitingInputStream is = new LimitingInputStream(new ByteArrayInputStream(TEST_BUFFER), 50);
- long bytesRead = StreamUtils.copy(is, new ByteArrayOutputStream());
+ long bytesRead = is.transferTo(new ByteArrayOutputStream());
assertEquals(bytesRead, TEST_BUFFER.length);
assertFalse(is.hasReachedLimit());
}
@@ -41,7 +41,7 @@ public void testReadLimitNotReached() throws IOException {
@Test
public void testReadLimitMatched() throws IOException {
final LimitingInputStream is = new LimitingInputStream(new ByteArrayInputStream(TEST_BUFFER), 10);
- long bytesRead = StreamUtils.copy(is, new ByteArrayOutputStream());
+ long bytesRead = is.transferTo(new ByteArrayOutputStream());
assertEquals(bytesRead, TEST_BUFFER.length);
assertTrue(is.hasReachedLimit());
}
@@ -49,7 +49,7 @@ public void testReadLimitMatched() throws IOException {
@Test
public void testReadLimitExceeded() throws IOException {
final LimitingInputStream is = new LimitingInputStream(new ByteArrayInputStream(TEST_BUFFER), 9);
- final long bytesRead = StreamUtils.copy(is, new ByteArrayOutputStream());
+ final long bytesRead = is.transferTo(new ByteArrayOutputStream());
assertEquals(9, bytesRead);
assertTrue(is.hasReachedLimit());
}
diff --git a/nifi-extension-bundles/nifi-atlassian-bundle/nifi-atlassian-extensions/src/main/java/org/apache/nifi/atlassian/bitbucket/BitbucketRepositoryClient.java b/nifi-extension-bundles/nifi-atlassian-bundle/nifi-atlassian-extensions/src/main/java/org/apache/nifi/atlassian/bitbucket/BitbucketRepositoryClient.java
index 5a62df207e08..354b8d1b6e8c 100644
--- a/nifi-extension-bundles/nifi-atlassian-bundle/nifi-atlassian-extensions/src/main/java/org/apache/nifi/atlassian/bitbucket/BitbucketRepositoryClient.java
+++ b/nifi-extension-bundles/nifi-atlassian-bundle/nifi-atlassian-extensions/src/main/java/org/apache/nifi/atlassian/bitbucket/BitbucketRepositoryClient.java
@@ -26,7 +26,6 @@
import org.apache.nifi.registry.flow.git.client.GitCommit;
import org.apache.nifi.registry.flow.git.client.GitCreateContentRequest;
import org.apache.nifi.registry.flow.git.client.GitRepositoryClient;
-import org.apache.nifi.stream.io.StreamUtils;
import org.apache.nifi.web.client.api.HttpResponseEntity;
import org.apache.nifi.web.client.api.HttpUriBuilder;
import org.apache.nifi.web.client.api.StandardHttpContentType;
@@ -1255,7 +1254,7 @@ private String getFileName(final String resolvedPath) {
private byte[] toByteArray(final InputStream inputStream) throws FlowRegistryException {
try (inputStream; ByteArrayOutputStream outputStream = new ByteArrayOutputStream()) {
- StreamUtils.copy(inputStream, outputStream);
+ inputStream.transferTo(outputStream);
return outputStream.toByteArray();
} catch (IOException e) {
throw new FlowRegistryException("Failed to prepare multipart request", e);
diff --git a/nifi-extension-bundles/nifi-compress-bundle/nifi-compress-processors/src/main/java/org/apache/nifi/processors/compress/ModifyCompression.java b/nifi-extension-bundles/nifi-compress-bundle/nifi-compress-processors/src/main/java/org/apache/nifi/processors/compress/ModifyCompression.java
index 31fd48101c6c..9355e9229bca 100644
--- a/nifi-extension-bundles/nifi-compress-bundle/nifi-compress-processors/src/main/java/org/apache/nifi/processors/compress/ModifyCompression.java
+++ b/nifi-extension-bundles/nifi-compress-bundle/nifi-compress-processors/src/main/java/org/apache/nifi/processors/compress/ModifyCompression.java
@@ -50,7 +50,6 @@
import org.apache.nifi.processors.compress.property.CompressionStrategy;
import org.apache.nifi.processors.compress.property.FilenameStrategy;
import org.apache.nifi.stream.io.GZIPOutputStream;
-import org.apache.nifi.stream.io.StreamUtils;
import org.apache.nifi.util.StopWatch;
import org.apache.nifi.util.StringUtils;
import org.tukaani.xz.LZMA2Options;
@@ -239,7 +238,7 @@ public void onTrigger(final ProcessContext context, final ProcessSession session
final BufferedOutputStream bufferedOutputStream = new BufferedOutputStream(flowFileOutputStream, STREAM_BUFFER_SIZE);
final OutputStream outputStream = getCompressionOutputStream(outputCompressionStrategy, outputCompressionLevel, mimeTypeRef, bufferedOutputStream)
) {
- StreamUtils.copy(inputStream, outputStream);
+ inputStream.transferTo(outputStream);
}
});
stopWatch.stop();
diff --git a/nifi-extension-bundles/nifi-email-bundle/nifi-email-processors/src/main/java/org/apache/nifi/processors/email/ExtractEmailAttachments.java b/nifi-extension-bundles/nifi-email-bundle/nifi-email-processors/src/main/java/org/apache/nifi/processors/email/ExtractEmailAttachments.java
index affb0cb05e0d..c2763df5c0cb 100644
--- a/nifi-extension-bundles/nifi-email-bundle/nifi-email-processors/src/main/java/org/apache/nifi/processors/email/ExtractEmailAttachments.java
+++ b/nifi-extension-bundles/nifi-email-bundle/nifi-email-processors/src/main/java/org/apache/nifi/processors/email/ExtractEmailAttachments.java
@@ -41,7 +41,6 @@
import org.apache.nifi.processor.ProcessSession;
import org.apache.nifi.processor.Relationship;
import org.apache.nifi.processor.exception.FlowFileHandlingException;
-import org.apache.nifi.stream.io.StreamUtils;
import java.io.BufferedInputStream;
import java.io.IOException;
@@ -135,7 +134,7 @@ public void onTrigger(final ProcessContext context, final ProcessSession session
String parentUuid = originalFlowFile.getAttribute(CoreAttributes.UUID.key());
attributes.put(ATTACHMENT_ORIGINAL_UUID, parentUuid);
attributes.put(ATTACHMENT_ORIGINAL_FILENAME, originalFlowFileName);
- split = session.append(split, out -> StreamUtils.copy(data.getInputStream(), out));
+ split = session.append(split, out -> data.getInputStream().transferTo(out));
split = session.putAllAttributes(split, attributes);
attachmentsList.add(split);
}
diff --git a/nifi-extension-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/PutHDFS.java b/nifi-extension-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/PutHDFS.java
index 4ab54d0129e2..35c21947031a 100644
--- a/nifi-extension-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/PutHDFS.java
+++ b/nifi-extension-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/PutHDFS.java
@@ -62,7 +62,6 @@
import org.apache.nifi.processor.util.StandardValidators;
import org.apache.nifi.processors.hadoop.util.GSSExceptionRollbackYieldSessionHandler;
import org.apache.nifi.processors.transfer.ResourceTransferSource;
-import org.apache.nifi.stream.io.StreamUtils;
import org.apache.nifi.util.StopWatch;
import java.io.BufferedInputStream;
@@ -444,7 +443,7 @@ public Object run() {
}
} else {
BufferedInputStream bis = new BufferedInputStream(in);
- StreamUtils.copy(bis, fos);
+ bis.transferTo(fos);
bis = null;
fos.flush();
}
diff --git a/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/main/java/org/apache/nifi/processors/pgp/DecryptContentPGP.java b/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/main/java/org/apache/nifi/processors/pgp/DecryptContentPGP.java
index a8016144cf0f..7252a3ab555e 100644
--- a/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/main/java/org/apache/nifi/processors/pgp/DecryptContentPGP.java
+++ b/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/main/java/org/apache/nifi/processors/pgp/DecryptContentPGP.java
@@ -42,7 +42,6 @@
import org.apache.nifi.processors.pgp.attributes.DecryptionStrategy;
import org.apache.nifi.processors.pgp.exception.PGPDecryptionException;
import org.apache.nifi.processors.pgp.exception.PGPProcessException;
-import org.apache.nifi.stream.io.StreamUtils;
import org.apache.nifi.util.StringUtils;
import org.bouncycastle.bcpg.KeyIdentifier;
import org.bouncycastle.openpgp.PGPCompressedData;
@@ -280,7 +279,7 @@ public void process(final InputStream inputStream, final OutputStream outputStre
if (DecryptionStrategy.PACKAGED == decryptionStrategy) {
try {
final InputStream decryptedDataStream = getDecryptedDataStream(encryptedData);
- StreamUtils.copy(decryptedDataStream, outputStream);
+ decryptedDataStream.transferTo(outputStream);
} catch (final PGPException e) {
final String message = String.format("PGP Decryption Failed [%s]", getEncryptedDataType(encryptedData));
throw new PGPDecryptionException(message, e);
@@ -291,7 +290,7 @@ public void process(final InputStream inputStream, final OutputStream outputStre
attributes.put(PGPAttributeKey.LITERAL_DATA_MODIFIED, Long.toString(literalData.getModificationTime().getTime()));
getLogger().debug("PGP Decrypted File Name [{}] Modified [{}]", literalData.getFileName(), literalData.getModificationTime());
- StreamUtils.copy(literalData.getInputStream(), outputStream);
+ literalData.getInputStream().transferTo(outputStream);
}
if (isVerified(encryptedData)) {
diff --git a/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/main/java/org/apache/nifi/processors/pgp/EncryptContentPGP.java b/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/main/java/org/apache/nifi/processors/pgp/EncryptContentPGP.java
index 98878a659da2..ac4cb7fee1a3 100644
--- a/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/main/java/org/apache/nifi/processors/pgp/EncryptContentPGP.java
+++ b/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/main/java/org/apache/nifi/processors/pgp/EncryptContentPGP.java
@@ -367,7 +367,7 @@ protected void processEncoding(final InputStream inputStream, final OutputStream
) {
if (isPacketFound(pushbackInputStream)) {
// Write OpenPGP packets to encrypted stream without additional encoding
- StreamUtils.copy(pushbackInputStream, encryptedOutputStream);
+ pushbackInputStream.transferTo(encryptedOutputStream);
} else {
super.processEncoding(pushbackInputStream, encryptedOutputStream);
}
diff --git a/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/main/java/org/apache/nifi/processors/pgp/VerifyContentPGP.java b/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/main/java/org/apache/nifi/processors/pgp/VerifyContentPGP.java
index c5d43c803f65..61532524e67f 100644
--- a/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/main/java/org/apache/nifi/processors/pgp/VerifyContentPGP.java
+++ b/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/main/java/org/apache/nifi/processors/pgp/VerifyContentPGP.java
@@ -35,7 +35,6 @@
import org.apache.nifi.processor.Relationship;
import org.apache.nifi.processor.io.StreamCallback;
import org.apache.nifi.processors.pgp.exception.PGPProcessException;
-import org.apache.nifi.stream.io.StreamUtils;
import org.bouncycastle.openpgp.PGPCompressedData;
import org.bouncycastle.openpgp.PGPException;
import org.bouncycastle.openpgp.PGPLiteralData;
@@ -260,7 +259,7 @@ private void processLiteralData(final PGPLiteralData literalData,
setLiteralDataAttributes(literalData);
final InputStream literalInputStream = literalData.getInputStream();
if (onePassSignature == null) {
- StreamUtils.copy(literalInputStream, outputStream);
+ literalInputStream.transferTo(outputStream);
} else {
processSignedStream(literalInputStream, outputStream, onePassSignature);
}
diff --git a/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/main/java/org/apache/nifi/processors/pgp/io/EncodingStreamCallback.java b/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/main/java/org/apache/nifi/processors/pgp/io/EncodingStreamCallback.java
index 68ef8e88ccad..7e0d11b0a126 100644
--- a/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/main/java/org/apache/nifi/processors/pgp/io/EncodingStreamCallback.java
+++ b/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/main/java/org/apache/nifi/processors/pgp/io/EncodingStreamCallback.java
@@ -20,7 +20,6 @@
import org.apache.nifi.processors.pgp.attributes.CompressionAlgorithm;
import org.apache.nifi.processors.pgp.attributes.FileEncoding;
import org.apache.nifi.processors.pgp.exception.PGPProcessException;
-import org.apache.nifi.stream.io.StreamUtils;
import org.bouncycastle.bcpg.ArmoredOutputStream;
import org.bouncycastle.openpgp.PGPCompressedDataGenerator;
import org.bouncycastle.openpgp.PGPException;
@@ -103,7 +102,7 @@ protected void processEncoding(final InputStream inputStream, final OutputStream
protected void processCompression(final InputStream inputStream, final OutputStream compressedOutputStream) throws IOException, PGPException {
final PGPLiteralDataGenerator generator = new PGPLiteralDataGenerator();
try (final OutputStream literalOutputStream = openLiteralOutputStream(generator, compressedOutputStream)) {
- StreamUtils.copy(inputStream, literalOutputStream);
+ inputStream.transferTo(literalOutputStream);
}
generator.close();
}
diff --git a/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/test/java/org/apache/nifi/processors/pgp/DecryptContentPGPTest.java b/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/test/java/org/apache/nifi/processors/pgp/DecryptContentPGPTest.java
index 9dde23ae33ba..31eaa7c55a7a 100644
--- a/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/test/java/org/apache/nifi/processors/pgp/DecryptContentPGPTest.java
+++ b/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/test/java/org/apache/nifi/processors/pgp/DecryptContentPGPTest.java
@@ -23,7 +23,6 @@
import org.apache.nifi.processors.pgp.exception.PGPDecryptionException;
import org.apache.nifi.processors.pgp.exception.PGPProcessException;
import org.apache.nifi.reporting.InitializationException;
-import org.apache.nifi.stream.io.StreamUtils;
import org.apache.nifi.util.LogMessage;
import org.apache.nifi.util.MockFlowFile;
import org.apache.nifi.util.TestRunner;
@@ -402,7 +401,7 @@ private void assertLiteralDataEquals(final Object object) throws IOException {
assertEquals(MODIFIED, literalData.getModificationTime());
final ByteArrayOutputStream outputStream = new ByteArrayOutputStream();
- StreamUtils.copy(literalData.getDataStream(), outputStream);
+ literalData.getDataStream().transferTo(outputStream);
final String literal = outputStream.toString(DATA_CHARSET);
assertEquals(DATA, literal);
}
diff --git a/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/test/java/org/apache/nifi/processors/pgp/EncryptContentPGPTest.java b/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/test/java/org/apache/nifi/processors/pgp/EncryptContentPGPTest.java
index a5987a280b1e..8b0ab82ce8f2 100644
--- a/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/test/java/org/apache/nifi/processors/pgp/EncryptContentPGPTest.java
+++ b/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/test/java/org/apache/nifi/processors/pgp/EncryptContentPGPTest.java
@@ -24,7 +24,6 @@
import org.apache.nifi.processors.pgp.attributes.FileEncoding;
import org.apache.nifi.processors.pgp.attributes.SymmetricKeyAlgorithm;
import org.apache.nifi.reporting.InitializationException;
-import org.apache.nifi.stream.io.StreamUtils;
import org.apache.nifi.util.MockFlowFile;
import org.apache.nifi.util.TestRunner;
import org.apache.nifi.util.TestRunners;
@@ -336,11 +335,11 @@ private byte[] getDecryptedData(final PGPPublicKeyEncryptedData publicKeyEncrypt
private byte[] getDecryptedData(final InputStream decryptedDataStream, final DecryptionStrategy decryptionStrategy) throws PGPException, IOException {
final ByteArrayOutputStream outputStream = new ByteArrayOutputStream();
if (DecryptionStrategy.PACKAGED == decryptionStrategy) {
- StreamUtils.copy(decryptedDataStream, outputStream);
+ decryptedDataStream.transferTo(outputStream);
} else {
final PGPObjectFactory objectFactory = new JcaPGPObjectFactory(decryptedDataStream);
final PGPLiteralData literalData = getLiteralData(objectFactory);
- StreamUtils.copy(literalData.getDataStream(), outputStream);
+ literalData.getDataStream().transferTo(outputStream);
}
return outputStream.toByteArray();
}
diff --git a/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/test/java/org/apache/nifi/processors/pgp/io/EncodingStreamCallbackTest.java b/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/test/java/org/apache/nifi/processors/pgp/io/EncodingStreamCallbackTest.java
index 3926a0bd3e85..d02e71b4ebaa 100644
--- a/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/test/java/org/apache/nifi/processors/pgp/io/EncodingStreamCallbackTest.java
+++ b/nifi-extension-bundles/nifi-pgp-bundle/nifi-pgp-processors/src/test/java/org/apache/nifi/processors/pgp/io/EncodingStreamCallbackTest.java
@@ -18,7 +18,6 @@
import org.apache.nifi.processors.pgp.attributes.CompressionAlgorithm;
import org.apache.nifi.processors.pgp.attributes.FileEncoding;
-import org.apache.nifi.stream.io.StreamUtils;
import org.bouncycastle.openpgp.PGPCompressedData;
import org.bouncycastle.openpgp.PGPException;
import org.bouncycastle.openpgp.PGPLiteralData;
@@ -90,7 +89,7 @@ private void assertLiteralDataEquals(final InputStream inputStream) throws IOExc
assertEquals(PGPLiteralData.BINARY, literalData.getFormat());
final ByteArrayOutputStream literalOutputStream = new ByteArrayOutputStream();
- StreamUtils.copy(literalData.getDataStream(), literalOutputStream);
+ literalData.getDataStream().transferTo(literalOutputStream);
assertArrayEquals(DATA, literalOutputStream.toByteArray());
}
}
diff --git a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/EncodeContent.java b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/EncodeContent.java
index 3dbfa12f30b4..dc0d4ae0eaf6 100644
--- a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/EncodeContent.java
+++ b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/EncodeContent.java
@@ -44,7 +44,6 @@
import org.apache.nifi.processors.standard.encoding.LineOutputMode;
import org.apache.nifi.processors.standard.util.ValidatingBase32InputStream;
import org.apache.nifi.processors.standard.util.ValidatingBase64InputStream;
-import org.apache.nifi.stream.io.StreamUtils;
import org.apache.nifi.util.StopWatch;
import java.io.IOException;
@@ -201,7 +200,7 @@ public void process(final InputStream in, final OutputStream out) throws IOExcep
.setLineSeparator(this.lineSeparator.getBytes())
.get())
.get()) {
- StreamUtils.copy(in, bos);
+ in.transferTo(bos);
}
}
}
@@ -211,7 +210,7 @@ private static class DecodeBase64 implements StreamCallback {
@Override
public void process(final InputStream in, final OutputStream out) throws IOException {
try (Base64InputStream bis = new Base64InputStream(new ValidatingBase64InputStream(in))) {
- StreamUtils.copy(bis, out);
+ bis.transferTo(out);
}
}
}
@@ -239,7 +238,7 @@ public void process(final InputStream in, final OutputStream out) throws IOExcep
.setLineSeparator(this.lineSeparator.getBytes())
.get())
.get()) {
- StreamUtils.copy(in, bos);
+ in.transferTo(bos);
}
}
}
@@ -249,7 +248,7 @@ private static class DecodeBase32 implements StreamCallback {
@Override
public void process(final InputStream in, final OutputStream out) throws IOException {
try (Base32InputStream bis = new Base32InputStream(new ValidatingBase32InputStream(in))) {
- StreamUtils.copy(bis, out);
+ bis.transferTo(out);
}
}
}
diff --git a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/ExecuteStreamCommand.java b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/ExecuteStreamCommand.java
index 11ddcf00bccb..f0ff0a3d9d49 100644
--- a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/ExecuteStreamCommand.java
+++ b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/ExecuteStreamCommand.java
@@ -51,7 +51,6 @@
import org.apache.nifi.processors.standard.util.ArgumentUtils;
import org.apache.nifi.processors.standard.util.SoftLimitBoundedByteArrayOutputStream;
import org.apache.nifi.stream.io.LimitingInputStream;
-import org.apache.nifi.stream.io.StreamUtils;
import java.io.BufferedInputStream;
import java.io.BufferedOutputStream;
@@ -569,7 +568,7 @@ public void process(final InputStream incomingFlowFileIS) throws IOException {
if (putToAttribute) {
try (SoftLimitBoundedByteArrayOutputStream softLimitBoundedBAOS = new SoftLimitBoundedByteArrayOutputStream(attributeSize)) {
readStdoutReadable(ignoreStdin, stdinWritable, logger, incomingFlowFileIS);
- final long longSize = StreamUtils.copy(stdoutReadable, softLimitBoundedBAOS);
+ final long longSize = stdoutReadable.transferTo(softLimitBoundedBAOS);
// Because the outputStream has a cap that the copy doesn't know about, adjust
// the actual size
@@ -591,7 +590,7 @@ public void process(final InputStream incomingFlowFileIS) throws IOException {
} else {
outputFlowFile = session.write(outputFlowFile, out -> {
readStdoutReadable(ignoreStdin, stdinWritable, logger, incomingFlowFileIS);
- StreamUtils.copy(stdoutReadable, out);
+ stdoutReadable.transferTo(out);
try {
exitCode = process.waitFor();
} catch (InterruptedException e) {
@@ -607,7 +606,7 @@ private static void readStdoutReadable(final boolean ignoreStdin, final OutputSt
Thread writerThread = new Thread(() -> {
if (!ignoreStdin) {
try {
- StreamUtils.copy(incomingFlowFileIS, stdinWritable);
+ incomingFlowFileIS.transferTo(stdinWritable);
} catch (IOException e) {
// This is unlikely to occur, and isn't handled at the moment
// Bug captured in NIFI-1194
diff --git a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/GetFileResource.java b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/GetFileResource.java
index 0c1446785227..d2d43102077a 100644
--- a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/GetFileResource.java
+++ b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/GetFileResource.java
@@ -37,7 +37,6 @@
import org.apache.nifi.processor.Relationship;
import org.apache.nifi.processor.util.StandardValidators;
import org.apache.nifi.scheduling.SchedulingStrategy;
-import org.apache.nifi.stream.io.StreamUtils;
import java.io.IOException;
import java.io.InputStream;
@@ -135,7 +134,7 @@ public void onTrigger(final ProcessContext context, final ProcessSession session
FlowFile flowFile = session.create();
try (final InputStream inputStream = context.getProperty(FILE_RESOURCE).asResource().read()) {
- flowFile = session.write(flowFile, out -> StreamUtils.copy(inputStream, out));
+ flowFile = session.write(flowFile, inputStream::transferTo);
} catch (IOException e) {
getLogger().error("Could not create FlowFile from Resource [{}]", context.getProperty(FILE_RESOURCE).getValue(), e);
}
diff --git a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/HandleHttpRequest.java b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/HandleHttpRequest.java
index 7e0279640114..816d0e940b39 100644
--- a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/HandleHttpRequest.java
+++ b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/HandleHttpRequest.java
@@ -64,7 +64,6 @@
import org.apache.nifi.processors.standard.util.HTTPUtils;
import org.apache.nifi.scheduling.ExecutionNode;
import org.apache.nifi.ssl.SSLContextProvider;
-import org.apache.nifi.stream.io.StreamUtils;
import org.eclipse.jetty.ee11.servlet.ServletContextHandler;
import org.eclipse.jetty.ee11.servlet.ServletContextRequest;
import org.eclipse.jetty.server.Connector;
@@ -658,7 +657,7 @@ public void onTrigger(final ProcessContext context, final ProcessSession session
Part part = parts.get(i);
FlowFile flowFile = session.create();
try (OutputStream flowFileOut = session.write(flowFile)) {
- StreamUtils.copy(part.getInputStream(), flowFileOut);
+ part.getInputStream().transferTo(flowFileOut);
} catch (IOException e) {
handleFlowContentStreamingError(session, container, Optional.of(flowFile), e);
return;
@@ -690,7 +689,7 @@ public void onTrigger(final ProcessContext context, final ProcessSession session
} else {
FlowFile flowFile = session.create();
try (OutputStream flowFileOut = session.write(flowFile)) {
- StreamUtils.copy(request.getInputStream(), flowFileOut);
+ request.getInputStream().transferTo(flowFileOut);
} catch (final IOException e) {
handleFlowContentStreamingError(session, container, Optional.of(flowFile), e);
return;
diff --git a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/MergeContent.java b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/MergeContent.java
index 10a86741701c..0d74bc8808d6 100644
--- a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/MergeContent.java
+++ b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/MergeContent.java
@@ -70,7 +70,6 @@
import org.apache.nifi.processors.standard.merge.AttributeStrategy;
import org.apache.nifi.processors.standard.merge.AttributeStrategyUtil;
import org.apache.nifi.stream.io.NonCloseableOutputStream;
-import org.apache.nifi.stream.io.StreamUtils;
import org.apache.nifi.util.FlowFilePackager;
import org.apache.nifi.util.FlowFilePackagerV1;
import org.apache.nifi.util.FlowFilePackagerV2;
@@ -807,7 +806,7 @@ public FlowFile merge(final Bin bin, final ProcessContext context) {
final Iterator itr = contents.iterator();
while (itr.hasNext()) {
final FlowFile flowFile = itr.next();
- bin.getSession().read(flowFile, in -> StreamUtils.copy(in, out));
+ bin.getSession().read(flowFile, in -> in.transferTo(out));
if (itr.hasNext()) {
if (demarcator != null) {
diff --git a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/TailFile.java b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/TailFile.java
index 137cde9c8a5e..3b4fc42e0e16 100644
--- a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/TailFile.java
+++ b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/TailFile.java
@@ -1488,7 +1488,7 @@ private TailFileState consumeFileFully(final File file, final ProcessContext con
try (final InputStream fis = new FileInputStream(file)) {
flowFile = session.write(flowFile, out -> {
flushLinesBuffer(out, new CRC32());
- StreamUtils.copy(fis, out);
+ fis.transferTo(out);
});
}
diff --git a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/UnpackContent.java b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/UnpackContent.java
index f1980640076a..3100e6982fb7 100644
--- a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/UnpackContent.java
+++ b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/UnpackContent.java
@@ -523,7 +523,7 @@ protected void processEntry(final InputStream zipInputStream, boolean directory,
attributes.put(FRAGMENT_ID, fragmentId);
attributes.put(FRAGMENT_INDEX, String.valueOf(++fragmentIndex));
unpackedFile = session.putAllAttributes(unpackedFile, attributes);
- unpackedFile = session.write(unpackedFile, outputStream -> StreamUtils.copy(zipInputStream, outputStream));
+ unpackedFile = session.write(unpackedFile, zipInputStream::transferTo);
} finally {
unpacked.add(unpackedFile);
}
diff --git a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/servlets/ListenHTTPServlet.java b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/servlets/ListenHTTPServlet.java
index 7145ffacce23..5c850d2649bf 100644
--- a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/servlets/ListenHTTPServlet.java
+++ b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/servlets/ListenHTTPServlet.java
@@ -49,7 +49,6 @@
import org.apache.nifi.serialization.RecordSetWriter;
import org.apache.nifi.serialization.RecordSetWriterFactory;
import org.apache.nifi.serialization.record.RecordSet;
-import org.apache.nifi.stream.io.StreamUtils;
import org.apache.nifi.util.FlowFileUnpackager;
import org.apache.nifi.util.FlowFileUnpackagerV1;
import org.apache.nifi.util.FlowFileUnpackagerV2;
@@ -297,7 +296,7 @@ private Set handleMultipartRequest(HttpServletRequest request, Process
OutputStream flowFileOutputStream = session.write(flowFile);
InputStream partInputStream = part.getInputStream()
) {
- StreamUtils.copy(partInputStream, flowFileOutputStream);
+ partInputStream.transferTo(flowFileOutputStream);
}
flowFile = saveRequestDetailsAsAttributes(request, session, foundSubject, foundIssuer, flowFile);
flowFile = savePartDetailsAsAttributes(session, part, flowFile, i, requestParts.size());
diff --git a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/util/FTPTransfer.java b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/util/FTPTransfer.java
index bbd8777754c1..973fd8f2c59a 100644
--- a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/util/FTPTransfer.java
+++ b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/util/FTPTransfer.java
@@ -39,7 +39,6 @@
import org.apache.nifi.processors.standard.ftp.StandardFTPClientProvider;
import org.apache.nifi.proxy.ProxyConfiguration;
import org.apache.nifi.proxy.ProxySpec;
-import org.apache.nifi.stream.io.StreamUtils;
import java.io.File;
import java.io.FileNotFoundException;
@@ -325,7 +324,7 @@ public FlowFile getRemoteFile(final String remoteFileName, final FlowFile origFl
throw new IOException(reply);
}
- resultFlowFile = session.write(origFlowFile, out -> StreamUtils.copy(in, out));
+ resultFlowFile = session.write(origFlowFile, in::transferTo);
client.completePendingCommand();
return resultFlowFile;
}
diff --git a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/util/SFTPTransfer.java b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/util/SFTPTransfer.java
index 3276d0f9b4c9..a618518692d1 100644
--- a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/util/SFTPTransfer.java
+++ b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/util/SFTPTransfer.java
@@ -38,7 +38,6 @@
import org.apache.nifi.processors.standard.ssh.StandardSshClientProvider;
import org.apache.nifi.proxy.ProxyConfiguration;
import org.apache.nifi.proxy.ProxySpec;
-import org.apache.nifi.stream.io.StreamUtils;
import org.apache.nifi.util.StringUtils;
import org.apache.sshd.client.session.ClientSession;
import org.apache.sshd.common.cipher.BuiltinCiphers;
@@ -532,7 +531,7 @@ public FlowFile getRemoteFile(final String remoteFileName, final FlowFile origFl
final SftpClient sftpClient = getSFTPClient(origFlowFile);
try (InputStream inputStream = sftpClient.read(remoteFileName)) {
- return session.write(origFlowFile, out -> StreamUtils.copy(inputStream, out));
+ return session.write(origFlowFile, inputStream::transferTo);
} catch (final SftpException e) {
final int status = e.getStatus();
switch (status) {
diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/StandardProcessSession.java b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/StandardProcessSession.java
index baa398ec8276..d7061d7d9b1f 100644
--- a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/StandardProcessSession.java
+++ b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/StandardProcessSession.java
@@ -3380,7 +3380,7 @@ public FlowFile append(FlowFile source, final OutputStreamCallback writer) {
appendableStreams.put(newClaim, outStream);
// We need to copy all of the data from the old claim to the new claim
- StreamUtils.copy(oldClaimIn, outStream);
+ oldClaimIn.transferTo(outStream);
// Don't allow flushing of the BufferedOutputStream. The callback may well call wrap our stream in another object that needs to be flushed.
// This is OK, but append() is often used many times to append just a small bit of data, over & over. If we allow flushing of our buffered output stream
diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/repository/FileSystemRepository.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/repository/FileSystemRepository.java
index 7994df080428..b8fba96fff85 100644
--- a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/repository/FileSystemRepository.java
+++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/repository/FileSystemRepository.java
@@ -771,7 +771,7 @@ public ContentClaim clone(final ContentClaim original, final boolean lossToleran
final ContentClaim newClaim = create(lossTolerant);
try (final InputStream in = read(original);
final OutputStream out = write(newClaim)) {
- StreamUtils.copy(in, out);
+ in.transferTo(out);
} catch (final IOException ioe) {
decrementClaimantCount(newClaim);
remove(newClaim);
@@ -790,7 +790,7 @@ public long importFrom(final Path content, final ContentClaim claim) throws IOEx
@Override
public long importFrom(final InputStream content, final ContentClaim claim) throws IOException {
try (final OutputStream out = write(claim, false)) {
- return StreamUtils.copy(content, out);
+ return content.transferTo(out);
}
}
@@ -806,7 +806,7 @@ public long exportTo(final ContentClaim claim, final Path destination, final boo
try (final InputStream in = read(claim);
final FileOutputStream fos = new FileOutputStream(destination.toFile(), append)) {
- final long copied = StreamUtils.copy(in, fos);
+ final long copied = in.transferTo(fos);
if (alwaysSync) {
fos.getFD().sync();
}
@@ -853,7 +853,7 @@ public long exportTo(final ContentClaim claim, final OutputStream destination) t
}
try (final InputStream in = read(claim)) {
- return StreamUtils.copy(in, destination);
+ return in.transferTo(destination);
}
}
diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/TestFileSystemSwapManager.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/TestFileSystemSwapManager.java
index 9a3c869daad5..4297048ecf51 100644
--- a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/TestFileSystemSwapManager.java
+++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/TestFileSystemSwapManager.java
@@ -25,7 +25,6 @@
import org.apache.nifi.controller.repository.claim.ResourceClaim;
import org.apache.nifi.controller.repository.claim.ResourceClaimManager;
import org.apache.nifi.events.EventReporter;
-import org.apache.nifi.stream.io.StreamUtils;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
@@ -120,7 +119,7 @@ public void testSwapFileUnknownToRepoNotSwappedIn() throws IOException {
final File originalSwapFile = new File("src/test/resources/swap/444-old-swap-file.swap");
try (final OutputStream fos = new FileOutputStream(targetFile);
final InputStream fis = new FileInputStream(originalSwapFile)) {
- StreamUtils.copy(fis, fos);
+ fis.transferTo(fos);
}
final FileSystemSwapManager swapManager = new FileSystemSwapManager(Paths.get("target"));
diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/repository/StandardProcessSessionIT.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/repository/StandardProcessSessionIT.java
index 96bb5b5f0cb8..582cd68d11f5 100644
--- a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/repository/StandardProcessSessionIT.java
+++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/repository/StandardProcessSessionIT.java
@@ -629,7 +629,7 @@ public void testWriteCallbackThenOutputStream() throws IOException {
private byte[] readContents(final FlowFile flowFile) throws IOException {
try (final InputStream in = session.read(flowFile);
final ByteArrayOutputStream baos = new ByteArrayOutputStream()) {
- StreamUtils.copy(in, baos);
+ in.transferTo(baos);
return baos.toByteArray();
}
}
@@ -1668,7 +1668,7 @@ public void testAppendToFlowFileWhereResourceClaimHasMultipleContentClaims() thr
// Read the content back and ensure that it is correct
final byte[] buff;
try (final ByteArrayOutputStream baos = new ByteArrayOutputStream()) {
- newSession.read(toUpdate, in -> StreamUtils.copy(in, baos));
+ newSession.read(toUpdate, in -> in.transferTo(baos));
buff = baos.toByteArray();
}
@@ -3368,7 +3368,7 @@ public long importFrom(Path content, ContentClaim claim) throws IOException {
public long importFrom(InputStream content, ContentClaim claim) throws IOException {
final long size;
try (final OutputStream out = write(claim)) {
- size = StreamUtils.copy(content, out);
+ size = content.transferTo(out);
}
((StandardContentClaim) claim).setLength(size);
return size;
@@ -3387,7 +3387,7 @@ public long exportTo(ContentClaim claim, Path destination, boolean append, long
@Override
public long exportTo(ContentClaim claim, OutputStream destination) throws IOException {
try (final InputStream in = read(claim)) {
- return StreamUtils.copy(in, destination);
+ return in.transferTo(destination);
}
}
diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/repository/TestFileSystemRepository.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/repository/TestFileSystemRepository.java
index b63842364997..ba1a1ba4b602 100644
--- a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/repository/TestFileSystemRepository.java
+++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/repository/TestFileSystemRepository.java
@@ -493,21 +493,21 @@ public void testWriteWithNoContent() throws IOException {
final ByteArrayOutputStream baos = new ByteArrayOutputStream();
try (final InputStream in = repository.read(claim1)) {
- StreamUtils.copy(in, baos);
+ in.transferTo(baos);
}
assertEquals("Hello", baos.toString());
baos.reset();
try (final InputStream in = repository.read(claim2)) {
- StreamUtils.copy(in, baos);
+ in.transferTo(baos);
}
assertEquals("", baos.toString());
assertEquals(0, baos.size());
baos.reset();
try (final InputStream in = repository.read(claim3)) {
- StreamUtils.copy(in, baos);
+ in.transferTo(baos);
}
assertEquals(" World", baos.toString());
}
@@ -576,7 +576,7 @@ public void testImportFromFile() throws IOException {
final ByteArrayOutputStream baos = new ByteArrayOutputStream();
try (final InputStream in = repository.read(claim)) {
- StreamUtils.copy(in, baos);
+ in.transferTo(baos);
}
assertArrayEquals(expected, baos.toByteArray());
diff --git a/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/extensions/DownloadQueue.java b/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/extensions/DownloadQueue.java
index 76e97b42eb0b..1dc8f14b511e 100644
--- a/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/extensions/DownloadQueue.java
+++ b/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/extensions/DownloadQueue.java
@@ -23,7 +23,6 @@
import org.apache.nifi.nar.ExtensionManager;
import org.apache.nifi.nar.NarClassLoaders;
import org.apache.nifi.nar.NarManifestEntry;
-import org.apache.nifi.stream.io.StreamUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -212,7 +211,7 @@ private File download(final BundleCoordinate coordinate) throws BundleNotFoundEx
// on startup, we could have two different threads attempting to download the same artifact. So we give the file a unique name by using the UUID.
final File tmpFile = new File(destinationFile.getParentFile(), destinationFile.getName() + ".download." + UUID.randomUUID());
try (final OutputStream out = new FileOutputStream(tmpFile)) {
- StreamUtils.copy(extensionStream, out);
+ extensionStream.transferTo(out);
}
// We need to rename our temporary file to the destination file. There's a chance that another thread could be finishing the same process
diff --git a/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/flow/StandardStatelessFlow.java b/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/flow/StandardStatelessFlow.java
index 5c317d7bf752..8671b8680506 100644
--- a/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/flow/StandardStatelessFlow.java
+++ b/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/flow/StandardStatelessFlow.java
@@ -72,7 +72,6 @@
import org.apache.nifi.stateless.repository.RepositoryContextFactory;
import org.apache.nifi.stateless.repository.StatelessProvenanceRepository;
import org.apache.nifi.stateless.session.AsynchronousCommitTracker;
-import org.apache.nifi.stream.io.StreamUtils;
import org.apache.nifi.util.Connectables;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -822,7 +821,7 @@ public QueueSize enqueue(final InputStream flowFileContents, final Map StreamUtils.copy(flowFileContents, out));
+ flowFile = session.write(flowFile, flowFileContents::transferTo);
flowFile = session.putAllAttributes(flowFile, attributes);
session.transfer(flowFile, LocalPort.PORT_RELATIONSHIP);
session.commitAsync();
diff --git a/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/repository/ByteArrayContentRepository.java b/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/repository/ByteArrayContentRepository.java
index 64aa56744dba..87d0cb5ac479 100644
--- a/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/repository/ByteArrayContentRepository.java
+++ b/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/repository/ByteArrayContentRepository.java
@@ -114,7 +114,7 @@ public ContentClaim clone(final ContentClaim original, final boolean lossToleran
final ContentClaim clone = create(lossTolerant);
try (final InputStream in = read(original);
final OutputStream out = write(clone)) {
- StreamUtils.copy(in, out);
+ in.transferTo(out);
}
return clone;
@@ -130,7 +130,7 @@ public long importFrom(final Path content, final ContentClaim claim) throws IOEx
@Override
public long importFrom(final InputStream content, final ContentClaim claim) throws IOException {
try (final OutputStream out = write(claim)) {
- return StreamUtils.copy(content, out);
+ return content.transferTo(out);
}
}
@@ -157,7 +157,7 @@ public long exportTo(final ContentClaim claim, final Path destination, final boo
@Override
public long exportTo(final ContentClaim claim, final OutputStream destination) throws IOException {
try (final InputStream in = read(claim)) {
- return StreamUtils.copy(in, destination);
+ return in.transferTo(destination);
}
}
diff --git a/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/repository/StatelessFileSystemContentRepository.java b/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/repository/StatelessFileSystemContentRepository.java
index 9d00a0345afb..3dbe36fb13e1 100644
--- a/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/repository/StatelessFileSystemContentRepository.java
+++ b/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/repository/StatelessFileSystemContentRepository.java
@@ -176,7 +176,7 @@ public ContentClaim clone(final ContentClaim original, final boolean lossToleran
final ContentClaim clone = create(lossTolerant);
try (final InputStream in = read(original);
final OutputStream out = write(clone)) {
- StreamUtils.copy(in, out);
+ in.transferTo(out);
}
return clone;
@@ -192,7 +192,7 @@ public long importFrom(final Path content, final ContentClaim claim) throws IOEx
@Override
public long importFrom(final InputStream content, final ContentClaim claim) throws IOException {
try (final OutputStream out = write(claim)) {
- return StreamUtils.copy(content, out);
+ return content.transferTo(out);
}
}
@@ -219,7 +219,7 @@ public long exportTo(final ContentClaim claim, final Path destination, final boo
@Override
public long exportTo(final ContentClaim claim, final OutputStream destination) throws IOException {
try (final InputStream in = read(claim)) {
- return StreamUtils.copy(in, destination);
+ return in.transferTo(destination);
}
}
diff --git a/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/test/java/org/apache/nifi/stateless/repository/TestStatelessFileSystemContentRepository.java b/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/test/java/org/apache/nifi/stateless/repository/TestStatelessFileSystemContentRepository.java
index 56787023f5df..0bdc9c4c1944 100644
--- a/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/test/java/org/apache/nifi/stateless/repository/TestStatelessFileSystemContentRepository.java
+++ b/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/test/java/org/apache/nifi/stateless/repository/TestStatelessFileSystemContentRepository.java
@@ -23,7 +23,6 @@
import org.apache.nifi.controller.repository.claim.ResourceClaimManager;
import org.apache.nifi.controller.repository.claim.StandardResourceClaimManager;
import org.apache.nifi.events.EventReporter;
-import org.apache.nifi.stream.io.StreamUtils;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -82,7 +81,7 @@ public void testWriteThenRead() throws IOException {
final byte[] bytesRead;
try (final InputStream in = repository.read(claim);
final ByteArrayOutputStream baos = new ByteArrayOutputStream()) {
- StreamUtils.copy(in, baos);
+ in.transferTo(baos);
bytesRead = baos.toByteArray();
}
diff --git a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/AssetReadingProcessor.java b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/AssetReadingProcessor.java
index 54053d1e8be3..bd7cf784ca3b 100644
--- a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/AssetReadingProcessor.java
+++ b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/AssetReadingProcessor.java
@@ -26,7 +26,6 @@
import org.apache.nifi.processor.Relationship;
import org.apache.nifi.processor.exception.ProcessException;
import org.apache.nifi.processor.util.StandardValidators;
-import org.apache.nifi.stream.io.StreamUtils;
import java.io.File;
import java.io.FileInputStream;
@@ -93,7 +92,7 @@ public void onTrigger(final ProcessContext context, final ProcessSession session
try {
try (final InputStream inputStream = new FileInputStream(sourceFile);
final OutputStream outputStream = Files.newOutputStream(temporaryPath)) {
- StreamUtils.copy(inputStream, outputStream);
+ inputStream.transferTo(outputStream);
}
try {
diff --git a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/ConcatenateFlowFiles.java b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/ConcatenateFlowFiles.java
index e66253381850..f7d3554d2297 100644
--- a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/ConcatenateFlowFiles.java
+++ b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/ConcatenateFlowFiles.java
@@ -29,7 +29,6 @@
import org.apache.nifi.processor.Relationship;
import org.apache.nifi.processor.exception.ProcessException;
import org.apache.nifi.processor.util.StandardValidators;
-import org.apache.nifi.stream.io.StreamUtils;
import java.io.InputStream;
import java.io.OutputStream;
@@ -103,7 +102,7 @@ public void onTrigger(final ProcessContext context, final ProcessSessionFactory
try (final OutputStream out = mergeSession.write(merged)) {
for (final FlowFile input : flowFiles) {
try (final InputStream in = mergeSession.read(input)) {
- StreamUtils.copy(in, out);
+ in.transferTo(out);
}
}
} catch (final Exception e) {
diff --git a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/UnzipFlowFile.java b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/UnzipFlowFile.java
index 6e16f1f50249..b8b6d75eac79 100644
--- a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/UnzipFlowFile.java
+++ b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/UnzipFlowFile.java
@@ -22,7 +22,6 @@
import org.apache.nifi.processor.ProcessSession;
import org.apache.nifi.processor.Relationship;
import org.apache.nifi.processor.exception.ProcessException;
-import org.apache.nifi.stream.io.StreamUtils;
import java.io.IOException;
import java.io.InputStream;
@@ -76,9 +75,7 @@ public void onTrigger(final ProcessContext context, final ProcessSession session
outFile = session.putAttribute(outFile, "filename", filename);
created.add(outFile);
- session.write(outFile, out -> {
- StreamUtils.copy(zipInputStream, out);
- });
+ session.write(outFile, zipInputStream::transferTo);
}
} catch (final IOException e) {
diff --git a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/VerifyContents.java b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/VerifyContents.java
index d54507d2c0ac..7b2ea14c5832 100644
--- a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/VerifyContents.java
+++ b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/VerifyContents.java
@@ -25,7 +25,6 @@
import org.apache.nifi.processor.ProcessSession;
import org.apache.nifi.processor.Relationship;
import org.apache.nifi.processor.exception.ProcessException;
-import org.apache.nifi.stream.io.StreamUtils;
import java.io.ByteArrayOutputStream;
import java.io.InputStream;
@@ -81,7 +80,7 @@ public void onTrigger(final ProcessContext context, final ProcessSession session
try (final InputStream in = session.read(flowFile);
final ByteArrayOutputStream baos = new ByteArrayOutputStream()) {
- StreamUtils.copy(in, baos);
+ in.transferTo(baos);
contents = baos.toString(StandardCharsets.UTF_8);
} catch (final Exception e) {
throw new ProcessException(e);
diff --git a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/WriteToFile.java b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/WriteToFile.java
index e5086cc9cd16..42053bf617ca 100644
--- a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/WriteToFile.java
+++ b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/WriteToFile.java
@@ -25,7 +25,6 @@
import org.apache.nifi.processor.ProcessSession;
import org.apache.nifi.processor.Relationship;
import org.apache.nifi.processor.exception.ProcessException;
-import org.apache.nifi.stream.io.StreamUtils;
import java.io.File;
import java.io.FileOutputStream;
@@ -88,7 +87,7 @@ public void onTrigger(final ProcessContext context, final ProcessSession session
try (final OutputStream out = new FileOutputStream(file);
final InputStream in = session.read(flowFile)) {
- StreamUtils.copy(in, out);
+ in.transferTo(out);
} catch (final Exception e) {
getLogger().error("Could not write FlowFile to {}", file.getAbsolutePath(), e);
session.transfer(flowFile, REL_FAILURE);
diff --git a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/NiFiClientUtil.java b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/NiFiClientUtil.java
index 2964c92106e8..faef5b9a5245 100644
--- a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/NiFiClientUtil.java
+++ b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/NiFiClientUtil.java
@@ -31,7 +31,6 @@
import org.apache.nifi.registry.flow.RegisteredFlowSnapshot;
import org.apache.nifi.remote.protocol.SiteToSiteTransportProtocol;
import org.apache.nifi.scheduling.ExecutionNode;
-import org.apache.nifi.stream.io.StreamUtils;
import org.apache.nifi.toolkit.client.ConnectionClient;
import org.apache.nifi.toolkit.client.ConnectorClient;
import org.apache.nifi.toolkit.client.NiFiClient;
@@ -2266,7 +2265,7 @@ public byte[] getFlowFileContentAsByteArray(final String connectionId, final int
try (final InputStream in = getFlowFileContent(connectionId, flowFileIndex);
final ByteArrayOutputStream baos = new ByteArrayOutputStream()) {
- StreamUtils.copy(in, baos);
+ in.transferTo(baos);
return baos.toByteArray();
}
}
diff --git a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/clustering/FlowSynchronizationIT.java b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/clustering/FlowSynchronizationIT.java
index 11fe85c43d19..d86454c2ff6b 100644
--- a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/clustering/FlowSynchronizationIT.java
+++ b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/clustering/FlowSynchronizationIT.java
@@ -23,7 +23,6 @@
import org.apache.nifi.controller.queue.LoadBalanceCompression;
import org.apache.nifi.controller.queue.LoadBalanceStrategy;
import org.apache.nifi.controller.service.ControllerServiceState;
-import org.apache.nifi.stream.io.StreamUtils;
import org.apache.nifi.tests.system.NiFiInstanceFactory;
import org.apache.nifi.tests.system.NiFiSystemIT;
import org.apache.nifi.toolkit.client.NiFiClientException;
@@ -819,7 +818,7 @@ private VersionedDataflow getNode2Flow() throws IOException {
try (final InputStream fis = new FileInputStream(flow);
final InputStream in = new GZIPInputStream(fis);
final ByteArrayOutputStream baos = new ByteArrayOutputStream()) {
- StreamUtils.copy(in, baos);
+ in.transferTo(baos);
bytes = baos.toByteArray();
}
diff --git a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/parameters/ClusteredProviderParamFlowSyncIT.java b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/parameters/ClusteredProviderParamFlowSyncIT.java
index e6d987fccc2c..c019fdbaf03a 100644
--- a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/parameters/ClusteredProviderParamFlowSyncIT.java
+++ b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/parameters/ClusteredProviderParamFlowSyncIT.java
@@ -22,7 +22,6 @@
import org.apache.nifi.flow.VersionedParameterContext;
import org.apache.nifi.parameter.ParameterSensitivity;
import org.apache.nifi.parameter.StandardParameterProviderConfiguration;
-import org.apache.nifi.stream.io.StreamUtils;
import org.apache.nifi.tests.system.NiFiInstance;
import org.apache.nifi.tests.system.NiFiInstanceFactory;
import org.apache.nifi.tests.system.NiFiSystemIT;
@@ -161,7 +160,7 @@ private void corruptProviderBackedParameter(final NiFiInstance node) throws IOEx
try (final InputStream fis = new FileInputStream(flowFile);
final InputStream in = new GZIPInputStream(fis);
final ByteArrayOutputStream baos = new ByteArrayOutputStream()) {
- StreamUtils.copy(in, baos);
+ in.transferTo(baos);
bytes = baos.toByteArray();
}
diff --git a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/restart/FlowFileRestorationIT.java b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/restart/FlowFileRestorationIT.java
index f98dd98f714a..cb9749ad8c0f 100644
--- a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/restart/FlowFileRestorationIT.java
+++ b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/restart/FlowFileRestorationIT.java
@@ -17,7 +17,6 @@
package org.apache.nifi.tests.system.restart;
-import org.apache.nifi.stream.io.StreamUtils;
import org.apache.nifi.tests.system.NiFiInstance;
import org.apache.nifi.tests.system.NiFiSystemIT;
import org.apache.nifi.toolkit.client.NiFiClientException;
@@ -90,7 +89,7 @@ private byte[] getFlowFileContents(final String connectionId, final int flowFile
try (final InputStream in = getClientUtil().getFlowFileContent(connectionId, flowFileIndex);
final ByteArrayOutputStream baos = new ByteArrayOutputStream()) {
- StreamUtils.copy(in, baos);
+ in.transferTo(baos);
return baos.toByteArray();
}
}