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