Repository: nifi Updated Branches: refs/heads/0.x f46ea411e -> af61bbeac
NIFI-2471 fix Hadoop configuration resources when talking to multiple Hadoop clusters This closes #779. Project: http://git-wip-us.apache.org/repos/asf/nifi/repo Commit: http://git-wip-us.apache.org/repos/asf/nifi/commit/af61bbea Tree: http://git-wip-us.apache.org/repos/asf/nifi/tree/af61bbea Diff: http://git-wip-us.apache.org/repos/asf/nifi/diff/af61bbea Branch: refs/heads/0.x Commit: af61bbeaced6038eafe836f4ee7911a01ebeab6b Parents: f46ea41 Author: Mike Moser <[email protected]> Authored: Wed Aug 3 15:58:57 2016 -0400 Committer: Mark Payne <[email protected]> Committed: Thu Aug 4 14:56:59 2016 -0400 ---------------------------------------------------------------------- .../nifi/processors/hadoop/AbstractHadoopProcessor.java | 8 +++++--- .../java/org/apache/nifi/processors/hadoop/ListHDFS.java | 1 + 2 files changed, 6 insertions(+), 3 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/nifi/blob/af61bbea/nifi-nar-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/AbstractHadoopProcessor.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/AbstractHadoopProcessor.java b/nifi-nar-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/AbstractHadoopProcessor.java index f96aa78..2058091 100644 --- a/nifi-nar-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/AbstractHadoopProcessor.java +++ b/nifi-nar-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/AbstractHadoopProcessor.java @@ -45,6 +45,7 @@ import org.apache.nifi.util.NiFiProperties; import org.apache.nifi.util.StringUtils; import javax.net.SocketFactory; + import java.io.File; import java.io.IOException; import java.net.InetSocketAddress; @@ -178,8 +179,8 @@ public abstract class AbstractHadoopProcessor extends AbstractProcessor { } /* - * If your subclass also has an @OnScheduled annotated method and you need hdfsResources in that method, then be sure to call super.abstractOnScheduled(context) - */ + * If your subclass also has an @OnScheduled annotated method and you need hdfsResources in that method, then be sure to call super.abstractOnScheduled(context) + */ @OnScheduled public final void abstractOnScheduled(ProcessContext context) throws IOException { try { @@ -255,6 +256,7 @@ public abstract class AbstractHadoopProcessor extends AbstractProcessor { // disable caching of Configuration and FileSystem objects, else we cannot reconfigure the processor without a complete // restart String disableCacheName = String.format("fs.%s.impl.disable.cache", FileSystem.getDefaultUri(config).getScheme()); + config.set(disableCacheName, "true"); // If kerberos is enabled, create the file system as the kerberos principal // -- use RESOURCE_LOCK to guarantee UserGroupInformation is accessed by only a single thread at at time @@ -274,7 +276,7 @@ public abstract class AbstractHadoopProcessor extends AbstractProcessor { fs = getFileSystemAsUser(config, ugi); } } - config.set(disableCacheName, "true"); + getLogger().info("Initialized a new HDFS File System with working dir: {} default block size: {} default replication: {} config: {}", new Object[] { fs.getWorkingDirectory(), fs.getDefaultBlockSize(new Path(dir)), fs.getDefaultReplication(new Path(dir)), config.toString() }); return new HdfsResources(config, fs, ugi); http://git-wip-us.apache.org/repos/asf/nifi/blob/af61bbea/nifi-nar-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/ListHDFS.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/ListHDFS.java b/nifi-nar-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/ListHDFS.java index f5daef2..b71e73d 100644 --- a/nifi-nar-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/ListHDFS.java +++ b/nifi-nar-bundles/nifi-hadoop-bundle/nifi-hdfs-processors/src/main/java/org/apache/nifi/processors/hadoop/ListHDFS.java @@ -163,6 +163,7 @@ public class ListHDFS extends AbstractHadoopProcessor { @Override public void onPropertyModified(final PropertyDescriptor descriptor, final String oldValue, final String newValue) { + super.onPropertyModified(descriptor, oldValue, newValue); if (isConfigurationRestored() && descriptor.equals(DIRECTORY)) { latestTimestampEmitted = -1L; latestTimestampListed = -1L;
