From b571e047f411970f58165df1789401e3e949c139 Mon Sep 17 00:00:00 2001 From: Shubham Sharma Date: Thu, 30 Jul 2026 15:39:18 -0400 Subject: [PATCH] NIFI-16159 Add support for FileResourceService in PutFile --- .../nifi-standard-processors/pom.xml | 15 +++++++ .../nifi/processors/standard/PutFile.java | 24 ++++++++++- .../nifi/processors/standard/TestPutFile.java | 41 +++++++++++++++++++ 3 files changed, 78 insertions(+), 2 deletions(-) diff --git a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/pom.xml b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/pom.xml index 6350480daebb..b9fad090088a 100644 --- a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/pom.xml +++ b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/pom.xml @@ -38,6 +38,21 @@ nifi-standard-record-utils 2.11.0-SNAPSHOT + + org.apache.nifi + nifi-resource-transfer + 2.11.0-SNAPSHOT + + + org.apache.nifi + nifi-file-resource-service-api + + + org.apache.nifi + nifi-file-resource-service + 2.11.0-SNAPSHOT + test + org.apache.commons commons-dbcp2 diff --git a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/PutFile.java b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/PutFile.java index db820f89ed3b..794f6fde5fdb 100644 --- a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/PutFile.java +++ b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/PutFile.java @@ -27,6 +27,7 @@ import org.apache.nifi.components.ValidationResult; import org.apache.nifi.components.Validator; import org.apache.nifi.expression.ExpressionLanguageScope; +import org.apache.nifi.fileresource.service.api.FileResource; import org.apache.nifi.flowfile.FlowFile; import org.apache.nifi.flowfile.attributes.CoreAttributes; import org.apache.nifi.logging.ComponentLog; @@ -36,11 +37,14 @@ import org.apache.nifi.processor.Relationship; import org.apache.nifi.processor.exception.ProcessException; import org.apache.nifi.processor.util.StandardValidators; +import org.apache.nifi.processors.transfer.ResourceTransferSource; import org.apache.nifi.util.StopWatch; +import java.io.InputStream; import java.nio.file.Files; import java.nio.file.Path; import java.nio.file.Paths; +import java.nio.file.StandardCopyOption; import java.nio.file.attribute.PosixFileAttributeView; import java.nio.file.attribute.PosixFilePermissions; import java.nio.file.attribute.UserPrincipalLookupService; @@ -48,11 +52,16 @@ import java.time.format.DateTimeFormatter; import java.util.Arrays; import java.util.List; +import java.util.Optional; import java.util.Set; import java.util.concurrent.TimeUnit; import java.util.regex.Matcher; import java.util.regex.Pattern; +import static org.apache.nifi.processors.transfer.ResourceTransferProperties.FILE_RESOURCE_SERVICE; +import static org.apache.nifi.processors.transfer.ResourceTransferProperties.RESOURCE_TRANSFER_SOURCE; +import static org.apache.nifi.processors.transfer.ResourceTransferUtils.getFileResource; + @SupportsBatching @InputRequirement(Requirement.INPUT_REQUIRED) @Tags({"put", "local", "copy", "archive", "files", "filesystem"}) @@ -158,7 +167,9 @@ public class PutFile extends AbstractProcessor { CHANGE_LAST_MODIFIED_TIME, CHANGE_PERMISSIONS, CHANGE_OWNER, - CHANGE_GROUP + CHANGE_GROUP, + RESOURCE_TRANSFER_SOURCE, + FILE_RESOURCE_SERVICE ); public static final int MAX_FILE_LOCK_ATTEMPTS = 10; @@ -304,7 +315,16 @@ public void onTrigger(final ProcessContext context, final ProcessSession session } } - session.exportTo(flowFile, dotCopyFile, false); + final ResourceTransferSource resourceTransferSource = ResourceTransferSource.valueOf( + context.getProperty(RESOURCE_TRANSFER_SOURCE).getValue()); + final Optional fileResource = getFileResource(resourceTransferSource, context, flowFile.getAttributes()); + if (fileResource.isPresent()) { + try (InputStream in = fileResource.get().getInputStream()) { + Files.copy(in, dotCopyFile, StandardCopyOption.REPLACE_EXISTING); + } + } else { + session.exportTo(flowFile, dotCopyFile, false); + } final String lastModifiedTime = context.getProperty(CHANGE_LAST_MODIFIED_TIME).evaluateAttributeExpressions(flowFile).getValue(); if (lastModifiedTime != null && !lastModifiedTime.isBlank()) { diff --git a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestPutFile.java b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestPutFile.java index 6afa8c97fe72..08ede2121038 100644 --- a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestPutFile.java +++ b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestPutFile.java @@ -16,7 +16,11 @@ */ package org.apache.nifi.processors.standard; +import org.apache.nifi.fileresource.service.StandardFileResourceService; +import org.apache.nifi.fileresource.service.api.FileResourceService; import org.apache.nifi.flowfile.attributes.CoreAttributes; +import org.apache.nifi.processors.transfer.ResourceTransferProperties; +import org.apache.nifi.processors.transfer.ResourceTransferSource; import org.apache.nifi.util.TestRunner; import org.apache.nifi.util.TestRunners; import org.junit.jupiter.api.AfterEach; @@ -27,6 +31,7 @@ import java.io.File; import java.io.IOException; +import java.nio.charset.StandardCharsets; import java.nio.file.FileVisitResult; import java.nio.file.FileVisitor; import java.nio.file.Files; @@ -244,6 +249,42 @@ public void testReplaceAndMaxFileLimitReach() throws IOException { assertEquals("Another file", new String(content)); } + @Test + public void testPutFileFromLocalFile() throws Exception { + final TestRunner runner = TestRunners.newTestRunner(new PutFile()); + runner.setProperty(PutFile.DIRECTORY, targetDir.getAbsolutePath()); + runner.setProperty(PutFile.CONFLICT_RESOLUTION, PutFile.REPLACE_RESOLUTION); + + final String attributeName = "file.path"; + final String serviceId = FileResourceService.class.getSimpleName(); + final FileResourceService service = new StandardFileResourceService(); + runner.addControllerService(serviceId, service); + runner.setProperty(service, StandardFileResourceService.FILE_PATH, String.format("${%s}", attributeName)); + runner.enableControllerService(service); + + runner.setProperty(ResourceTransferProperties.RESOURCE_TRANSFER_SOURCE, ResourceTransferSource.FILE_RESOURCE_SERVICE.getValue()); + runner.setProperty(ResourceTransferProperties.FILE_RESOURCE_SERVICE, serviceId); + + final byte[] fileData = "0123456789".getBytes(StandardCharsets.UTF_8); + final Path tempFilePath = Files.createTempFile("PutFile_testPutFileFromLocalFile_", ""); + Files.write(tempFilePath, fileData); + + try { + final Map attributes = new HashMap<>(); + attributes.put(CoreAttributes.FILENAME.key(), "targetFile.txt"); + attributes.put(attributeName, tempFilePath.toString()); + runner.enqueue(new byte[0], attributes); + runner.run(); + + runner.assertAllFlowFilesTransferred(PutFile.REL_SUCCESS, 1); + final Path targetPath = Paths.get(TARGET_DIRECTORY + "/targetFile.txt"); + final byte[] content = Files.readAllBytes(targetPath); + assertEquals("0123456789", new String(content, StandardCharsets.UTF_8)); + } finally { + Files.deleteIfExists(tempFilePath); + } + } + private TestRunner putFileRunner; private final String testFile = "src" + File.separator + "test" + File.separator + "resources" + File.separator + "hello.txt";