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";