From 4e00a330fe4c0a036ff4f4fec831a4458cd8ac08 Mon Sep 17 00:00:00 2001 From: dan-s1 Date: Fri, 7 Aug 2026 18:01:10 +0000 Subject: [PATCH 1/2] NIFI-16162 Deprecated methods copy and skip in org.apache.nifi.stream.io.StreamUtils and replaced them with java.io.Inputstream method equivalents --- .../nifi/remote/codec/StandardFlowFileCodec.java | 2 +- .../nifi/remote/util/SiteToSiteRestApiClient.java | 3 +-- .../nifi/remote/protocol/SiteToSiteTestUtils.java | 3 +-- .../org/apache/nifi/stream/io/StreamUtils.java | 11 +++++++++++ .../nifi/stream/io/LimitingInputStreamTest.java | 6 +++--- .../bitbucket/BitbucketRepositoryClient.java | 3 +-- .../processors/compress/ModifyCompression.java | 3 +-- .../processors/email/ExtractEmailAttachments.java | 3 +-- .../org/apache/nifi/processors/hadoop/PutHDFS.java | 3 +-- .../parquet/shared/NifiSeekableInputStream.java | 3 +-- .../nifi/processors/pgp/DecryptContentPGP.java | 5 ++--- .../nifi/processors/pgp/EncryptContentPGP.java | 2 +- .../nifi/processors/pgp/VerifyContentPGP.java | 3 +-- .../processors/pgp/io/EncodingStreamCallback.java | 3 +-- .../nifi/processors/pgp/DecryptContentPGPTest.java | 3 +-- .../nifi/processors/pgp/EncryptContentPGPTest.java | 5 ++--- .../pgp/io/EncodingStreamCallbackTest.java | 3 +-- .../nifi/processors/standard/EncodeContent.java | 9 ++++----- .../processors/standard/ExecuteStreamCommand.java | 7 +++---- .../nifi/processors/standard/GetFileResource.java | 3 +-- .../processors/standard/HandleHttpRequest.java | 5 ++--- .../nifi/processors/standard/MergeContent.java | 3 +-- .../apache/nifi/processors/standard/TailFile.java | 2 +- .../nifi/processors/standard/UnpackContent.java | 2 +- .../standard/servlets/ListenHTTPServlet.java | 3 +-- .../nifi/processors/standard/util/FTPTransfer.java | 3 +-- .../processors/standard/util/SFTPTransfer.java | 3 +-- .../provenance/EventIdFirstSchemaRecordReader.java | 2 +- .../serialization/CompressableRecordReader.java | 7 +++---- .../repository/StandardProcessSession.java | 6 +++--- .../repository/io/ContentClaimInputStream.java | 5 ++--- .../org/apache/nifi/controller/FlowController.java | 3 +-- .../clustered/ContentRepositoryFlowFileAccess.java | 3 +-- .../repository/FileSystemRepository.java | 14 +++++++------- .../nifi/controller/TestFileSystemSwapManager.java | 3 +-- .../repository/StandardProcessSessionIT.java | 8 ++++---- .../repository/TestFileSystemRepository.java | 8 ++++---- .../org/apache/nifi/extensions/DownloadQueue.java | 3 +-- .../engine/StandardExecutionProgress.java | 2 +- .../nifi/stateless/flow/StandardStatelessFlow.java | 3 +-- .../repository/ByteArrayContentRepository.java | 8 ++++---- .../StatelessFileSystemContentRepository.java | 10 +++++----- .../TestStatelessFileSystemContentRepository.java | 3 +-- .../tests/system/AssetReadingProcessor.java | 3 +-- .../tests/system/ConcatenateFlowFiles.java | 3 +-- .../processors/tests/system/UnzipFlowFile.java | 5 +---- .../processors/tests/system/VerifyContents.java | 3 +-- .../nifi/processors/tests/system/WriteToFile.java | 3 +-- .../apache/nifi/tests/system/NiFiClientUtil.java | 3 +-- .../system/clustering/FlowSynchronizationIT.java | 3 +-- .../ClusteredProviderParamFlowSyncIT.java | 3 +-- .../system/restart/FlowFileRestorationIT.java | 3 +-- 52 files changed, 96 insertions(+), 125 deletions(-) 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..6321acf7ff42 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; @@ -239,7 +248,9 @@ public static byte[] copyExclusive(final InputStream in, final OutputStream out, * @param stream the stream to skip over * @param bytesToSkip the number of bytes to skip * @throws IOException if any issues reading or skipping underlying stream + * @deprecated Use {@link InputStream#skipNBytes(long)} instead. */ + @Deprecated(since = "2.12.0", forRemoval = true) public static void skip(final InputStream stream, final long bytesToSkip) throws IOException { if (bytesToSkip <= 0) { return; 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-parquet-bundle/nifi-parquet-shared/src/main/java/org/apache/nifi/parquet/shared/NifiSeekableInputStream.java b/nifi-extension-bundles/nifi-parquet-bundle/nifi-parquet-shared/src/main/java/org/apache/nifi/parquet/shared/NifiSeekableInputStream.java index f3e348be51cf..42bbe29f334f 100644 --- a/nifi-extension-bundles/nifi-parquet-bundle/nifi-parquet-shared/src/main/java/org/apache/nifi/parquet/shared/NifiSeekableInputStream.java +++ b/nifi-extension-bundles/nifi-parquet-bundle/nifi-parquet-shared/src/main/java/org/apache/nifi/parquet/shared/NifiSeekableInputStream.java @@ -17,7 +17,6 @@ package org.apache.nifi.parquet.shared; import org.apache.nifi.stream.io.ByteCountingInputStream; -import org.apache.nifi.stream.io.StreamUtils; import org.apache.parquet.io.DelegatingSeekableInputStream; import java.io.IOException; @@ -51,7 +50,7 @@ public void seek(long newPos) throws IOException { } // must call getPos() again in case reset was called above - StreamUtils.skip(input, newPos - getPos()); + input.skipNBytes(newPos - getPos()); } @Override 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-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/EventIdFirstSchemaRecordReader.java b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/EventIdFirstSchemaRecordReader.java index b8be66760907..270be7943f4f 100644 --- a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/EventIdFirstSchemaRecordReader.java +++ b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/EventIdFirstSchemaRecordReader.java @@ -170,7 +170,7 @@ protected Optional readToEvent(final long eventId return Optional.ofNullable(event); } else { // This is not the record we want. Skip over it instead of deserializing it. - StreamUtils.skip(dis, recordLength); + dis.skipNBytes(recordLength); } } diff --git a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/serialization/CompressableRecordReader.java b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/serialization/CompressableRecordReader.java index 18801b991c63..a32989dfd41d 100644 --- a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/serialization/CompressableRecordReader.java +++ b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/serialization/CompressableRecordReader.java @@ -22,7 +22,6 @@ import org.apache.nifi.provenance.toc.TocReader; import org.apache.nifi.stream.io.ByteCountingInputStream; import org.apache.nifi.stream.io.LimitingInputStream; -import org.apache.nifi.stream.io.StreamUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -119,7 +118,7 @@ public void skipToBlock(final int blockIndex) throws IOException { final long bytesToSkip = offset - curOffset; if (bytesToSkip >= 0) { try { - StreamUtils.skip(rawInputStream, bytesToSkip); + rawInputStream.skipNBytes(bytesToSkip); logger.debug("Skipped stream from offset {} to {} ({} bytes skipped)", curOffset, offset, bytesToSkip); } catch (final EOFException eof) { throw new EOFException("Attempted to skip to byte offset " + offset + " for " + filename + " but file does not have that many bytes (TOC Reader=" + getTocReader() + ")"); @@ -245,7 +244,7 @@ public void close() throws IOException { @Override public void skip(final long bytesToSkip) throws IOException { - StreamUtils.skip(dis, bytesToSkip); + dis.skipNBytes(bytesToSkip); } @Override @@ -262,7 +261,7 @@ public void skipTo(final long position) throws IOException { } final long toSkip = position - currentPosition; - StreamUtils.skip(dis, toSkip); + dis.skipNBytes(toSkip); } protected String getFilename() { 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..8cfd8f14c3f8 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 @@ -2776,7 +2776,7 @@ private InputStream getInputStream(final FlowFile flowFile, final ContentClaim c if (currentReadClaimStream != null && currentReadClaimStream.getBytesConsumed() <= resourceClaimOffset) { final long bytesToSkip = resourceClaimOffset - currentReadClaimStream.getBytesConsumed(); if (bytesToSkip > 0) { - StreamUtils.skip(currentReadClaimStream, bytesToSkip); + currentReadClaimStream.skipNBytes(bytesToSkip); } final InputStream limitingInputStream = new LimitingInputStream(new DisableOnCloseInputStream(currentReadClaimStream), flowFile.getSize()); @@ -2799,7 +2799,7 @@ private InputStream getInputStream(final FlowFile flowFile, final ContentClaim c performanceTracker.endContentRead(); final InputStream performanceTrackInputStream = new PerformanceTrackingInputStream(contentRepoStream, performanceTracker); - StreamUtils.skip(performanceTrackInputStream, claim.getOffset() + contentClaimOffset); + performanceTrackInputStream.skipNBytes(claim.getOffset() + contentClaimOffset); final InputStream bufferedContentStream = new BufferedInputStream(contentRepoStream); final ByteCountingInputStream byteCountingInputStream = new ByteCountingInputStream(bufferedContentStream, claim.getOffset() + contentClaimOffset); currentReadClaimStream = byteCountingInputStream; @@ -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-components/src/main/java/org/apache/nifi/controller/repository/io/ContentClaimInputStream.java b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/io/ContentClaimInputStream.java index bfb9f5d651f4..3c129890b9d5 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/io/ContentClaimInputStream.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/io/ContentClaimInputStream.java @@ -20,7 +20,6 @@ import org.apache.nifi.controller.repository.claim.ContentClaim; import org.apache.nifi.controller.repository.metrics.PerformanceTracker; import org.apache.nifi.controller.repository.metrics.PerformanceTrackingInputStream; -import org.apache.nifi.stream.io.StreamUtils; import java.io.BufferedInputStream; import java.io.IOException; @@ -215,7 +214,7 @@ public void reset() throws IOException { performanceTracker.beginContentRead(); try { - StreamUtils.skip(delegate, markOffset - claimOffset); + delegate.skipNBytes(markOffset - claimOffset); } finally { performanceTracker.endContentRead(); } @@ -243,7 +242,7 @@ private void formDelegate() throws IOException { performanceTracker.beginContentRead(); try { delegate = new PerformanceTrackingInputStream(contentRepository.read(contentClaim), performanceTracker); - StreamUtils.skip(delegate, claimOffset); + delegate.skipNBytes(claimOffset); currentOffset = claimOffset; } finally { performanceTracker.endContentRead(); diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/FlowController.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/FlowController.java index 4953965ff887..e4225312578a 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/FlowController.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/FlowController.java @@ -226,7 +226,6 @@ import org.apache.nifi.scheduling.SchedulingStrategy; import org.apache.nifi.services.FlowService; import org.apache.nifi.stream.io.LimitingInputStream; -import org.apache.nifi.stream.io.StreamUtils; import org.apache.nifi.util.ComponentIdGenerator; import org.apache.nifi.util.FormatUtils; import org.apache.nifi.util.NiFiProperties; @@ -3404,7 +3403,7 @@ public InputStream getContent(final FlowFileRecord flowFile, final String reques stream = contentRepository.read(flowFile.getContentClaim()); final long contentClaimOffset = flowFile.getContentClaimOffset(); if (contentClaimOffset > 0L) { - StreamUtils.skip(stream, contentClaimOffset); + stream.skipNBytes(contentClaimOffset); } stream = new LimitingInputStream(stream, flowFile.getSize()); diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/queue/clustered/ContentRepositoryFlowFileAccess.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/queue/clustered/ContentRepositoryFlowFileAccess.java index 28ae05d92a79..6f1bba912f35 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/queue/clustered/ContentRepositoryFlowFileAccess.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/queue/clustered/ContentRepositoryFlowFileAccess.java @@ -21,7 +21,6 @@ import org.apache.nifi.controller.repository.ContentRepository; import org.apache.nifi.controller.repository.FlowFileRecord; import org.apache.nifi.controller.repository.io.LimitedInputStream; -import org.apache.nifi.stream.io.StreamUtils; import java.io.EOFException; import java.io.FilterInputStream; @@ -46,7 +45,7 @@ public InputStream read(final FlowFileRecord flowFile) throws IOException { if (flowFile.getContentClaimOffset() > 0) { try { - StreamUtils.skip(rawIn, flowFile.getContentClaimOffset()); + rawIn.skipNBytes(flowFile.getContentClaimOffset()); } catch (final EOFException eof) { throw new ContentNotFoundException(flowFile, flowFile.getContentClaim(), "FlowFile has a Content Claim Offset of " + flowFile.getContentClaimOffset() + " bytes but the Content Claim does not have that many bytes"); 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..4f9ac619ad94 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(); } @@ -836,7 +836,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)) { if (offset > 0) { - StreamUtils.skip(in, offset); + in.skipNBytes(offset); } StreamUtils.copy(in, fos, length); if (alwaysSync) { @@ -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); } } @@ -870,7 +870,7 @@ public long exportTo(final ContentClaim claim, final OutputStream destination, f return exportTo(claim, destination); } try (final InputStream in = read(claim)) { - StreamUtils.skip(in, offset); + in.skipNBytes(offset); final byte[] buffer = new byte[8192]; int len; long copied = 0L; @@ -924,7 +924,7 @@ public InputStream read(final ContentClaim claim) throws IOException { final InputStream fis = getInputStream(claim); if (claim.getOffset() > 0L) { try { - StreamUtils.skip(fis, claim.getOffset()); + fis.skipNBytes(claim.getOffset()); } catch (final EOFException eof) { closeQuietly(fis); 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/engine/StandardExecutionProgress.java b/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/engine/StandardExecutionProgress.java index 10a237831d46..8be41c16c261 100644 --- a/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/engine/StandardExecutionProgress.java +++ b/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/engine/StandardExecutionProgress.java @@ -255,7 +255,7 @@ public InputStream readContent(final FlowFile flowFile) throws IOException { final InputStream in = contentRepository.read(contentClaim); final long offset = flowFileRecord.getContentClaimOffset(); if (offset > 0) { - StreamUtils.skip(in, offset); + in.skipNBytes(offset); } return new LimitedInputStream(in, flowFile.getSize()); 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..767180a00427 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,14 +157,14 @@ 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); } } @Override public long exportTo(final ContentClaim claim, final OutputStream destination, final long offset, final long length) throws IOException { try (final InputStream in = read(claim)) { - StreamUtils.skip(in, offset); + in.skipNBytes(offset); StreamUtils.copy(in, destination, length); } 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..bbc32f6c1f4f 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,14 +219,14 @@ 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); } } @Override public long exportTo(final ContentClaim claim, final OutputStream destination, final long offset, final long length) throws IOException { try (final InputStream in = read(claim)) { - StreamUtils.skip(in, offset); + in.skipNBytes(offset); StreamUtils.copy(in, destination, length); } @@ -250,7 +250,7 @@ public InputStream read(final ContentClaim claim) throws IOException { } final InputStream resourceClaimIn = read(claim.getResourceClaim()); - StreamUtils.skip(resourceClaimIn, claim.getOffset()); + resourceClaimIn.skipNBytes(claim.getOffset()); final InputStream limitedIn = new LimitedInputStream(resourceClaimIn, claim.getLength()); return limitedIn; 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(); } } From 8222b0a0a101267ce505b2f5566bb468087cea87 Mon Sep 17 00:00:00 2001 From: dan-s1 Date: Fri, 7 Aug 2026 21:27:49 +0000 Subject: [PATCH 2/2] NIFI-16162 Backed out of changes from skip to skipNBytes --- .../main/java/org/apache/nifi/stream/io/StreamUtils.java | 2 -- .../nifi/parquet/shared/NifiSeekableInputStream.java | 3 ++- .../nifi/provenance/EventIdFirstSchemaRecordReader.java | 2 +- .../provenance/serialization/CompressableRecordReader.java | 7 ++++--- .../nifi/controller/repository/StandardProcessSession.java | 4 ++-- .../controller/repository/io/ContentClaimInputStream.java | 5 +++-- .../java/org/apache/nifi/controller/FlowController.java | 3 ++- .../queue/clustered/ContentRepositoryFlowFileAccess.java | 3 ++- .../nifi/controller/repository/FileSystemRepository.java | 6 +++--- .../nifi/stateless/engine/StandardExecutionProgress.java | 2 +- .../stateless/repository/ByteArrayContentRepository.java | 2 +- .../repository/StatelessFileSystemContentRepository.java | 4 ++-- 12 files changed, 23 insertions(+), 20 deletions(-) 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 6321acf7ff42..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 @@ -248,9 +248,7 @@ public static byte[] copyExclusive(final InputStream in, final OutputStream out, * @param stream the stream to skip over * @param bytesToSkip the number of bytes to skip * @throws IOException if any issues reading or skipping underlying stream - * @deprecated Use {@link InputStream#skipNBytes(long)} instead. */ - @Deprecated(since = "2.12.0", forRemoval = true) public static void skip(final InputStream stream, final long bytesToSkip) throws IOException { if (bytesToSkip <= 0) { return; diff --git a/nifi-extension-bundles/nifi-parquet-bundle/nifi-parquet-shared/src/main/java/org/apache/nifi/parquet/shared/NifiSeekableInputStream.java b/nifi-extension-bundles/nifi-parquet-bundle/nifi-parquet-shared/src/main/java/org/apache/nifi/parquet/shared/NifiSeekableInputStream.java index 42bbe29f334f..f3e348be51cf 100644 --- a/nifi-extension-bundles/nifi-parquet-bundle/nifi-parquet-shared/src/main/java/org/apache/nifi/parquet/shared/NifiSeekableInputStream.java +++ b/nifi-extension-bundles/nifi-parquet-bundle/nifi-parquet-shared/src/main/java/org/apache/nifi/parquet/shared/NifiSeekableInputStream.java @@ -17,6 +17,7 @@ package org.apache.nifi.parquet.shared; import org.apache.nifi.stream.io.ByteCountingInputStream; +import org.apache.nifi.stream.io.StreamUtils; import org.apache.parquet.io.DelegatingSeekableInputStream; import java.io.IOException; @@ -50,7 +51,7 @@ public void seek(long newPos) throws IOException { } // must call getPos() again in case reset was called above - input.skipNBytes(newPos - getPos()); + StreamUtils.skip(input, newPos - getPos()); } @Override diff --git a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/EventIdFirstSchemaRecordReader.java b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/EventIdFirstSchemaRecordReader.java index 270be7943f4f..b8be66760907 100644 --- a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/EventIdFirstSchemaRecordReader.java +++ b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/EventIdFirstSchemaRecordReader.java @@ -170,7 +170,7 @@ protected Optional readToEvent(final long eventId return Optional.ofNullable(event); } else { // This is not the record we want. Skip over it instead of deserializing it. - dis.skipNBytes(recordLength); + StreamUtils.skip(dis, recordLength); } } diff --git a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/serialization/CompressableRecordReader.java b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/serialization/CompressableRecordReader.java index a32989dfd41d..18801b991c63 100644 --- a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/serialization/CompressableRecordReader.java +++ b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/serialization/CompressableRecordReader.java @@ -22,6 +22,7 @@ import org.apache.nifi.provenance.toc.TocReader; import org.apache.nifi.stream.io.ByteCountingInputStream; import org.apache.nifi.stream.io.LimitingInputStream; +import org.apache.nifi.stream.io.StreamUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -118,7 +119,7 @@ public void skipToBlock(final int blockIndex) throws IOException { final long bytesToSkip = offset - curOffset; if (bytesToSkip >= 0) { try { - rawInputStream.skipNBytes(bytesToSkip); + StreamUtils.skip(rawInputStream, bytesToSkip); logger.debug("Skipped stream from offset {} to {} ({} bytes skipped)", curOffset, offset, bytesToSkip); } catch (final EOFException eof) { throw new EOFException("Attempted to skip to byte offset " + offset + " for " + filename + " but file does not have that many bytes (TOC Reader=" + getTocReader() + ")"); @@ -244,7 +245,7 @@ public void close() throws IOException { @Override public void skip(final long bytesToSkip) throws IOException { - dis.skipNBytes(bytesToSkip); + StreamUtils.skip(dis, bytesToSkip); } @Override @@ -261,7 +262,7 @@ public void skipTo(final long position) throws IOException { } final long toSkip = position - currentPosition; - dis.skipNBytes(toSkip); + StreamUtils.skip(dis, toSkip); } protected String getFilename() { 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 8cfd8f14c3f8..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 @@ -2776,7 +2776,7 @@ private InputStream getInputStream(final FlowFile flowFile, final ContentClaim c if (currentReadClaimStream != null && currentReadClaimStream.getBytesConsumed() <= resourceClaimOffset) { final long bytesToSkip = resourceClaimOffset - currentReadClaimStream.getBytesConsumed(); if (bytesToSkip > 0) { - currentReadClaimStream.skipNBytes(bytesToSkip); + StreamUtils.skip(currentReadClaimStream, bytesToSkip); } final InputStream limitingInputStream = new LimitingInputStream(new DisableOnCloseInputStream(currentReadClaimStream), flowFile.getSize()); @@ -2799,7 +2799,7 @@ private InputStream getInputStream(final FlowFile flowFile, final ContentClaim c performanceTracker.endContentRead(); final InputStream performanceTrackInputStream = new PerformanceTrackingInputStream(contentRepoStream, performanceTracker); - performanceTrackInputStream.skipNBytes(claim.getOffset() + contentClaimOffset); + StreamUtils.skip(performanceTrackInputStream, claim.getOffset() + contentClaimOffset); final InputStream bufferedContentStream = new BufferedInputStream(contentRepoStream); final ByteCountingInputStream byteCountingInputStream = new ByteCountingInputStream(bufferedContentStream, claim.getOffset() + contentClaimOffset); currentReadClaimStream = byteCountingInputStream; diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/io/ContentClaimInputStream.java b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/io/ContentClaimInputStream.java index 3c129890b9d5..bfb9f5d651f4 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/io/ContentClaimInputStream.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/repository/io/ContentClaimInputStream.java @@ -20,6 +20,7 @@ import org.apache.nifi.controller.repository.claim.ContentClaim; import org.apache.nifi.controller.repository.metrics.PerformanceTracker; import org.apache.nifi.controller.repository.metrics.PerformanceTrackingInputStream; +import org.apache.nifi.stream.io.StreamUtils; import java.io.BufferedInputStream; import java.io.IOException; @@ -214,7 +215,7 @@ public void reset() throws IOException { performanceTracker.beginContentRead(); try { - delegate.skipNBytes(markOffset - claimOffset); + StreamUtils.skip(delegate, markOffset - claimOffset); } finally { performanceTracker.endContentRead(); } @@ -242,7 +243,7 @@ private void formDelegate() throws IOException { performanceTracker.beginContentRead(); try { delegate = new PerformanceTrackingInputStream(contentRepository.read(contentClaim), performanceTracker); - delegate.skipNBytes(claimOffset); + StreamUtils.skip(delegate, claimOffset); currentOffset = claimOffset; } finally { performanceTracker.endContentRead(); diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/FlowController.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/FlowController.java index e4225312578a..4953965ff887 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/FlowController.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/FlowController.java @@ -226,6 +226,7 @@ import org.apache.nifi.scheduling.SchedulingStrategy; import org.apache.nifi.services.FlowService; import org.apache.nifi.stream.io.LimitingInputStream; +import org.apache.nifi.stream.io.StreamUtils; import org.apache.nifi.util.ComponentIdGenerator; import org.apache.nifi.util.FormatUtils; import org.apache.nifi.util.NiFiProperties; @@ -3403,7 +3404,7 @@ public InputStream getContent(final FlowFileRecord flowFile, final String reques stream = contentRepository.read(flowFile.getContentClaim()); final long contentClaimOffset = flowFile.getContentClaimOffset(); if (contentClaimOffset > 0L) { - stream.skipNBytes(contentClaimOffset); + StreamUtils.skip(stream, contentClaimOffset); } stream = new LimitingInputStream(stream, flowFile.getSize()); diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/queue/clustered/ContentRepositoryFlowFileAccess.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/queue/clustered/ContentRepositoryFlowFileAccess.java index 6f1bba912f35..28ae05d92a79 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/queue/clustered/ContentRepositoryFlowFileAccess.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/queue/clustered/ContentRepositoryFlowFileAccess.java @@ -21,6 +21,7 @@ import org.apache.nifi.controller.repository.ContentRepository; import org.apache.nifi.controller.repository.FlowFileRecord; import org.apache.nifi.controller.repository.io.LimitedInputStream; +import org.apache.nifi.stream.io.StreamUtils; import java.io.EOFException; import java.io.FilterInputStream; @@ -45,7 +46,7 @@ public InputStream read(final FlowFileRecord flowFile) throws IOException { if (flowFile.getContentClaimOffset() > 0) { try { - rawIn.skipNBytes(flowFile.getContentClaimOffset()); + StreamUtils.skip(rawIn, flowFile.getContentClaimOffset()); } catch (final EOFException eof) { throw new ContentNotFoundException(flowFile, flowFile.getContentClaim(), "FlowFile has a Content Claim Offset of " + flowFile.getContentClaimOffset() + " bytes but the Content Claim does not have that many bytes"); 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 4f9ac619ad94..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 @@ -836,7 +836,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)) { if (offset > 0) { - in.skipNBytes(offset); + StreamUtils.skip(in, offset); } StreamUtils.copy(in, fos, length); if (alwaysSync) { @@ -870,7 +870,7 @@ public long exportTo(final ContentClaim claim, final OutputStream destination, f return exportTo(claim, destination); } try (final InputStream in = read(claim)) { - in.skipNBytes(offset); + StreamUtils.skip(in, offset); final byte[] buffer = new byte[8192]; int len; long copied = 0L; @@ -924,7 +924,7 @@ public InputStream read(final ContentClaim claim) throws IOException { final InputStream fis = getInputStream(claim); if (claim.getOffset() > 0L) { try { - fis.skipNBytes(claim.getOffset()); + StreamUtils.skip(fis, claim.getOffset()); } catch (final EOFException eof) { closeQuietly(fis); diff --git a/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/engine/StandardExecutionProgress.java b/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/engine/StandardExecutionProgress.java index 8be41c16c261..10a237831d46 100644 --- a/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/engine/StandardExecutionProgress.java +++ b/nifi-stateless/nifi-stateless-bundle/nifi-stateless-engine/src/main/java/org/apache/nifi/stateless/engine/StandardExecutionProgress.java @@ -255,7 +255,7 @@ public InputStream readContent(final FlowFile flowFile) throws IOException { final InputStream in = contentRepository.read(contentClaim); final long offset = flowFileRecord.getContentClaimOffset(); if (offset > 0) { - in.skipNBytes(offset); + StreamUtils.skip(in, offset); } return new LimitedInputStream(in, flowFile.getSize()); 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 767180a00427..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 @@ -164,7 +164,7 @@ public long exportTo(final ContentClaim claim, final OutputStream destination) t @Override public long exportTo(final ContentClaim claim, final OutputStream destination, final long offset, final long length) throws IOException { try (final InputStream in = read(claim)) { - in.skipNBytes(offset); + StreamUtils.skip(in, offset); StreamUtils.copy(in, destination, length); } 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 bbc32f6c1f4f..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 @@ -226,7 +226,7 @@ public long exportTo(final ContentClaim claim, final OutputStream destination) t @Override public long exportTo(final ContentClaim claim, final OutputStream destination, final long offset, final long length) throws IOException { try (final InputStream in = read(claim)) { - in.skipNBytes(offset); + StreamUtils.skip(in, offset); StreamUtils.copy(in, destination, length); } @@ -250,7 +250,7 @@ public InputStream read(final ContentClaim claim) throws IOException { } final InputStream resourceClaimIn = read(claim.getResourceClaim()); - resourceClaimIn.skipNBytes(claim.getOffset()); + StreamUtils.skip(resourceClaimIn, claim.getOffset()); final InputStream limitedIn = new LimitedInputStream(resourceClaimIn, claim.getLength()); return limitedIn;