Repository: nifi
Updated Branches:
  refs/heads/master 868808f4b -> 4544f3969


NIFI-5152: MoveHDFS now works even with no upstream connection

This closes #2681.

Signed-off-by: Bryan Bende <[email protected]>


Project: http://git-wip-us.apache.org/repos/asf/nifi/repo
Commit: http://git-wip-us.apache.org/repos/asf/nifi/commit/4544f396
Tree: http://git-wip-us.apache.org/repos/asf/nifi/tree/4544f396
Diff: http://git-wip-us.apache.org/repos/asf/nifi/diff/4544f396

Branch: refs/heads/master
Commit: 4544f3969d6978f6637c59206c978281d377d1c3
Parents: 868808f
Author: zenfenan <[email protected]>
Authored: Sun May 6 09:46:19 2018 +0530
Committer: Bryan Bende <[email protected]>
Committed: Mon May 7 13:20:56 2018 -0400

----------------------------------------------------------------------
 .../apache/nifi/processors/hadoop/MoveHDFS.java | 35 ++++++++++----------
 1 file changed, 17 insertions(+), 18 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/nifi/blob/4544f396/nifi-nar-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/MoveHDFS.java
----------------------------------------------------------------------
diff --git 
a/nifi-nar-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/MoveHDFS.java
 
b/nifi-nar-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/MoveHDFS.java
index ea61ed1..a2292f5 100644
--- 
a/nifi-nar-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/MoveHDFS.java
+++ 
b/nifi-nar-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/MoveHDFS.java
@@ -24,6 +24,8 @@ import org.apache.hadoop.fs.FileUtil;
 import org.apache.hadoop.fs.Path;
 import org.apache.hadoop.fs.PathFilter;
 import org.apache.hadoop.security.UserGroupInformation;
+import org.apache.nifi.annotation.behavior.InputRequirement;
+import org.apache.nifi.annotation.behavior.InputRequirement.Requirement;
 import org.apache.nifi.annotation.behavior.ReadsAttribute;
 import org.apache.nifi.annotation.behavior.Restricted;
 import org.apache.nifi.annotation.behavior.Restriction;
@@ -65,6 +67,7 @@ import java.util.regex.Pattern;
 /**
  * This processor renames files on HDFS.
  */
+@InputRequirement(Requirement.INPUT_ALLOWED)
 @Tags({"hadoop", "HDFS", "put", "move", "filesystem", "moveHDFS"})
 @CapabilityDescription("Rename existing files or a directory of files 
(non-recursive) on Hadoop Distributed File System (HDFS).")
 @ReadsAttribute(attribute = "filename", description = "The name of the file 
written to HDFS comes from the value of this attribute.")
@@ -139,6 +142,7 @@ public class MoveHDFS extends AbstractHadoopProcessor {
             .name("Input Directory or File")
             .description("The HDFS directory from which files should be read, 
or a single file to read.")
             .defaultValue("${path}")
+            .required(true)
             
.addValidator(StandardValidators.ATTRIBUTE_EXPRESSION_LANGUAGE_VALIDATOR)
             
.expressionLanguageSupported(ExpressionLanguageScope.FLOWFILE_ATTRIBUTES)
             .build();
@@ -228,14 +232,16 @@ public class MoveHDFS extends AbstractHadoopProcessor {
 
     @Override
     public void onTrigger(ProcessContext context, ProcessSession session) 
throws ProcessException {
-        // MoveHDFS
-        FlowFile parentFlowFile = session.get();
-        if (parentFlowFile == null) {
+        FlowFile flowFile = session.get();
+
+        if (flowFile == null && context.hasIncomingConnection()) {
             return;
         }
 
+        flowFile = (flowFile != null) ? flowFile : session.create();
+
         final FileSystem hdfs = getFileSystem();
-        final String filenameValue = 
context.getProperty(INPUT_DIRECTORY_OR_FILE).evaluateAttributeExpressions(parentFlowFile).getValue();
+        final String filenameValue = 
context.getProperty(INPUT_DIRECTORY_OR_FILE).evaluateAttributeExpressions(flowFile).getValue();
 
         Path inputPath = null;
         try {
@@ -244,10 +250,10 @@ public class MoveHDFS extends AbstractHadoopProcessor {
                 throw new IOException("Input Directory or File does not exist 
in HDFS");
             }
         } catch (Exception e) {
-            getLogger().error("Failed to retrieve content from {} for {} due 
to {}; routing to failure", new Object[]{filenameValue, parentFlowFile, e});
-            parentFlowFile = session.putAttribute(parentFlowFile, 
"hdfs.failure.reason", e.getMessage());
-            parentFlowFile = session.penalize(parentFlowFile);
-            session.transfer(parentFlowFile, REL_FAILURE);
+            getLogger().error("Failed to retrieve content from {} for {} due 
to {}; routing to failure", new Object[]{filenameValue, flowFile, e});
+            flowFile = session.putAttribute(flowFile, "hdfs.failure.reason", 
e.getMessage());
+            flowFile = session.penalize(flowFile);
+            session.transfer(flowFile, REL_FAILURE);
             return;
         }
 
@@ -298,7 +304,7 @@ public class MoveHDFS extends AbstractHadoopProcessor {
             filePathQueue.drainTo(files);
             if (files.isEmpty()) {
                 // nothing to do!
-                session.remove(parentFlowFile);
+                session.remove(flowFile);
                 context.yield();
                 return;
             }
@@ -306,7 +312,7 @@ public class MoveHDFS extends AbstractHadoopProcessor {
             queueLock.unlock();
         }
 
-        processBatchOfFiles(files, context, session, parentFlowFile);
+        processBatchOfFiles(files, context, session, flowFile);
 
         queueLock.lock();
         try {
@@ -315,7 +321,7 @@ public class MoveHDFS extends AbstractHadoopProcessor {
             queueLock.unlock();
         }
 
-        session.remove(parentFlowFile);
+        session.remove(flowFile);
     }
 
     protected void processBatchOfFiles(final List<Path> files, final 
ProcessContext context,
@@ -496,7 +502,6 @@ public class MoveHDFS extends AbstractHadoopProcessor {
 
         final private String conflictResolution;
         final private String operation;
-        final private Path inputRootDirPath;
         final private Path outputRootDirPath;
         final private Pattern fileFilterPattern;
         final private boolean ignoreDottedFiles;
@@ -504,8 +509,6 @@ public class MoveHDFS extends AbstractHadoopProcessor {
         ProcessorConfiguration(final ProcessContext context) {
             conflictResolution = 
context.getProperty(CONFLICT_RESOLUTION).getValue();
             operation = context.getProperty(OPERATION).getValue();
-            final String inputDirValue = 
context.getProperty(INPUT_DIRECTORY_OR_FILE).evaluateAttributeExpressions().getValue();
-            inputRootDirPath = new Path(inputDirValue);
             final String outputDirValue = 
context.getProperty(OUTPUT_DIRECTORY).evaluateAttributeExpressions().getValue();
             outputRootDirPath = new Path(outputDirValue);
             final String fileFilterRegex = 
context.getProperty(FILE_FILTER_REGEX).getValue();
@@ -521,10 +524,6 @@ public class MoveHDFS extends AbstractHadoopProcessor {
             return conflictResolution;
         }
 
-        public Path getInput() {
-            return inputRootDirPath;
-        }
-
         public Path getOutputDirectory() {
             return outputRootDirPath;
         }

Reply via email to