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(); } }