YARN-4002. Make ResourceTrackerService#nodeHeartbeat more concurrent. Contributed by Rohith Sharma K S & Zhiguo Hong
This commit is contained in:
parent
141873ca7d
commit
feb90ffcca
@ -21,6 +21,9 @@
|
|||||||
import java.io.*;
|
import java.io.*;
|
||||||
import java.util.Set;
|
import java.util.Set;
|
||||||
import java.util.HashSet;
|
import java.util.HashSet;
|
||||||
|
import java.util.concurrent.locks.ReentrantReadWriteLock;
|
||||||
|
import java.util.concurrent.locks.ReentrantReadWriteLock.ReadLock;
|
||||||
|
import java.util.concurrent.locks.ReentrantReadWriteLock.WriteLock;
|
||||||
|
|
||||||
import org.apache.commons.io.Charsets;
|
import org.apache.commons.io.Charsets;
|
||||||
import org.apache.commons.logging.LogFactory;
|
import org.apache.commons.logging.LogFactory;
|
||||||
@ -38,6 +41,8 @@ public class HostsFileReader {
|
|||||||
private Set<String> excludes;
|
private Set<String> excludes;
|
||||||
private String includesFile;
|
private String includesFile;
|
||||||
private String excludesFile;
|
private String excludesFile;
|
||||||
|
private WriteLock writeLock;
|
||||||
|
private ReadLock readLock;
|
||||||
|
|
||||||
private static final Log LOG = LogFactory.getLog(HostsFileReader.class);
|
private static final Log LOG = LogFactory.getLog(HostsFileReader.class);
|
||||||
|
|
||||||
@ -47,6 +52,9 @@ public HostsFileReader(String inFile,
|
|||||||
excludes = new HashSet<String>();
|
excludes = new HashSet<String>();
|
||||||
includesFile = inFile;
|
includesFile = inFile;
|
||||||
excludesFile = exFile;
|
excludesFile = exFile;
|
||||||
|
ReentrantReadWriteLock rwLock = new ReentrantReadWriteLock();
|
||||||
|
this.writeLock = rwLock.writeLock();
|
||||||
|
this.readLock = rwLock.readLock();
|
||||||
refresh();
|
refresh();
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -57,6 +65,9 @@ public HostsFileReader(String includesFile, InputStream inFileInputStream,
|
|||||||
excludes = new HashSet<String>();
|
excludes = new HashSet<String>();
|
||||||
this.includesFile = includesFile;
|
this.includesFile = includesFile;
|
||||||
this.excludesFile = excludesFile;
|
this.excludesFile = excludesFile;
|
||||||
|
ReentrantReadWriteLock rwLock = new ReentrantReadWriteLock();
|
||||||
|
this.writeLock = rwLock.writeLock();
|
||||||
|
this.readLock = rwLock.readLock();
|
||||||
refresh(inFileInputStream, exFileInputStream);
|
refresh(inFileInputStream, exFileInputStream);
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -101,80 +112,126 @@ public static void readFileToSetWithFileInputStream(String type,
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
public synchronized void refresh() throws IOException {
|
public void refresh() throws IOException {
|
||||||
LOG.info("Refreshing hosts (include/exclude) list");
|
this.writeLock.lock();
|
||||||
Set<String> newIncludes = new HashSet<String>();
|
try {
|
||||||
Set<String> newExcludes = new HashSet<String>();
|
refresh(includesFile, excludesFile);
|
||||||
boolean switchIncludes = false;
|
} finally {
|
||||||
boolean switchExcludes = false;
|
this.writeLock.unlock();
|
||||||
if (!includesFile.isEmpty()) {
|
|
||||||
readFileToSet("included", includesFile, newIncludes);
|
|
||||||
switchIncludes = true;
|
|
||||||
}
|
|
||||||
if (!excludesFile.isEmpty()) {
|
|
||||||
readFileToSet("excluded", excludesFile, newExcludes);
|
|
||||||
switchExcludes = true;
|
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
if (switchIncludes) {
|
public void refresh(String includeFiles, String excludeFiles)
|
||||||
// switch the new hosts that are to be included
|
throws IOException {
|
||||||
includes = newIncludes;
|
LOG.info("Refreshing hosts (include/exclude) list");
|
||||||
}
|
this.writeLock.lock();
|
||||||
if (switchExcludes) {
|
try {
|
||||||
// switch the excluded hosts
|
// update instance variables
|
||||||
excludes = newExcludes;
|
updateFileNames(includeFiles, excludeFiles);
|
||||||
|
Set<String> newIncludes = new HashSet<String>();
|
||||||
|
Set<String> newExcludes = new HashSet<String>();
|
||||||
|
boolean switchIncludes = false;
|
||||||
|
boolean switchExcludes = false;
|
||||||
|
if (includeFiles != null && !includeFiles.isEmpty()) {
|
||||||
|
readFileToSet("included", includeFiles, newIncludes);
|
||||||
|
switchIncludes = true;
|
||||||
|
}
|
||||||
|
if (excludeFiles != null && !excludeFiles.isEmpty()) {
|
||||||
|
readFileToSet("excluded", excludeFiles, newExcludes);
|
||||||
|
switchExcludes = true;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (switchIncludes) {
|
||||||
|
// switch the new hosts that are to be included
|
||||||
|
includes = newIncludes;
|
||||||
|
}
|
||||||
|
if (switchExcludes) {
|
||||||
|
// switch the excluded hosts
|
||||||
|
excludes = newExcludes;
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
|
this.writeLock.unlock();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@Private
|
@Private
|
||||||
public synchronized void refresh(InputStream inFileInputStream,
|
public void refresh(InputStream inFileInputStream,
|
||||||
InputStream exFileInputStream) throws IOException {
|
InputStream exFileInputStream) throws IOException {
|
||||||
LOG.info("Refreshing hosts (include/exclude) list");
|
LOG.info("Refreshing hosts (include/exclude) list");
|
||||||
Set<String> newIncludes = new HashSet<String>();
|
this.writeLock.lock();
|
||||||
Set<String> newExcludes = new HashSet<String>();
|
try {
|
||||||
boolean switchIncludes = false;
|
Set<String> newIncludes = new HashSet<String>();
|
||||||
boolean switchExcludes = false;
|
Set<String> newExcludes = new HashSet<String>();
|
||||||
if (inFileInputStream != null) {
|
boolean switchIncludes = false;
|
||||||
readFileToSetWithFileInputStream("included", includesFile,
|
boolean switchExcludes = false;
|
||||||
inFileInputStream, newIncludes);
|
if (inFileInputStream != null) {
|
||||||
switchIncludes = true;
|
readFileToSetWithFileInputStream("included", includesFile,
|
||||||
}
|
inFileInputStream, newIncludes);
|
||||||
if (exFileInputStream != null) {
|
switchIncludes = true;
|
||||||
readFileToSetWithFileInputStream("excluded", excludesFile,
|
}
|
||||||
exFileInputStream, newExcludes);
|
if (exFileInputStream != null) {
|
||||||
switchExcludes = true;
|
readFileToSetWithFileInputStream("excluded", excludesFile,
|
||||||
}
|
exFileInputStream, newExcludes);
|
||||||
if (switchIncludes) {
|
switchExcludes = true;
|
||||||
// switch the new hosts that are to be included
|
}
|
||||||
includes = newIncludes;
|
if (switchIncludes) {
|
||||||
}
|
// switch the new hosts that are to be included
|
||||||
if (switchExcludes) {
|
includes = newIncludes;
|
||||||
// switch the excluded hosts
|
}
|
||||||
excludes = newExcludes;
|
if (switchExcludes) {
|
||||||
|
// switch the excluded hosts
|
||||||
|
excludes = newExcludes;
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
|
this.writeLock.unlock();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
public synchronized Set<String> getHosts() {
|
public Set<String> getHosts() {
|
||||||
return includes;
|
this.readLock.lock();
|
||||||
|
try {
|
||||||
|
return includes;
|
||||||
|
} finally {
|
||||||
|
this.readLock.unlock();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
public synchronized Set<String> getExcludedHosts() {
|
public Set<String> getExcludedHosts() {
|
||||||
return excludes;
|
this.readLock.lock();
|
||||||
|
try {
|
||||||
|
return excludes;
|
||||||
|
} finally {
|
||||||
|
this.readLock.unlock();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
public synchronized void setIncludesFile(String includesFile) {
|
public void getHostDetails(Set<String> includes, Set<String> excludes) {
|
||||||
|
this.readLock.lock();
|
||||||
|
try {
|
||||||
|
includes.addAll(this.includes);
|
||||||
|
excludes.addAll(this.excludes);
|
||||||
|
} finally {
|
||||||
|
this.readLock.unlock();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setIncludesFile(String includesFile) {
|
||||||
LOG.info("Setting the includes file to " + includesFile);
|
LOG.info("Setting the includes file to " + includesFile);
|
||||||
this.includesFile = includesFile;
|
this.includesFile = includesFile;
|
||||||
}
|
}
|
||||||
|
|
||||||
public synchronized void setExcludesFile(String excludesFile) {
|
public void setExcludesFile(String excludesFile) {
|
||||||
LOG.info("Setting the excludes file to " + excludesFile);
|
LOG.info("Setting the excludes file to " + excludesFile);
|
||||||
this.excludesFile = excludesFile;
|
this.excludesFile = excludesFile;
|
||||||
}
|
}
|
||||||
|
|
||||||
public synchronized void updateFileNames(String includesFile,
|
public void updateFileNames(String includeFiles, String excludeFiles) {
|
||||||
String excludesFile) {
|
this.writeLock.lock();
|
||||||
setIncludesFile(includesFile);
|
try {
|
||||||
setExcludesFile(excludesFile);
|
setIncludesFile(includeFiles);
|
||||||
|
setExcludesFile(excludeFiles);
|
||||||
|
} finally {
|
||||||
|
this.writeLock.unlock();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
@ -20,6 +20,8 @@
|
|||||||
import java.io.File;
|
import java.io.File;
|
||||||
import java.io.FileNotFoundException;
|
import java.io.FileNotFoundException;
|
||||||
import java.io.FileWriter;
|
import java.io.FileWriter;
|
||||||
|
import java.util.HashSet;
|
||||||
|
import java.util.Set;
|
||||||
|
|
||||||
import org.apache.hadoop.test.GenericTestUtils;
|
import org.apache.hadoop.test.GenericTestUtils;
|
||||||
import org.junit.*;
|
import org.junit.*;
|
||||||
@ -96,6 +98,30 @@ public void testHostsFileReader() throws Exception {
|
|||||||
assertTrue(hfp.getExcludedHosts().contains("somehost5"));
|
assertTrue(hfp.getExcludedHosts().contains("somehost5"));
|
||||||
assertFalse(hfp.getExcludedHosts().contains("host4"));
|
assertFalse(hfp.getExcludedHosts().contains("host4"));
|
||||||
|
|
||||||
|
// test for refreshing hostreader wit new include/exclude host files
|
||||||
|
String newExcludesFile = HOSTS_TEST_DIR + "/dfs1.exclude";
|
||||||
|
String newIncludesFile = HOSTS_TEST_DIR + "/dfs1.include";
|
||||||
|
|
||||||
|
efw = new FileWriter(newExcludesFile);
|
||||||
|
ifw = new FileWriter(newIncludesFile);
|
||||||
|
|
||||||
|
efw.write("#DFS-Hosts-excluded\n");
|
||||||
|
efw.write("node1\n");
|
||||||
|
efw.close();
|
||||||
|
|
||||||
|
ifw.write("#Hosts-in-DFS\n");
|
||||||
|
ifw.write("node2\n");
|
||||||
|
ifw.close();
|
||||||
|
|
||||||
|
hfp.refresh(newIncludesFile, newExcludesFile);
|
||||||
|
assertTrue(hfp.getExcludedHosts().contains("node1"));
|
||||||
|
assertTrue(hfp.getHosts().contains("node2"));
|
||||||
|
|
||||||
|
Set<String> hostsList = new HashSet<String>();
|
||||||
|
Set<String> excludeList = new HashSet<String>();
|
||||||
|
hfp.getHostDetails(hostsList, excludeList);
|
||||||
|
assertTrue(excludeList.contains("node1"));
|
||||||
|
assertTrue(hostsList.contains("node2"));
|
||||||
}
|
}
|
||||||
|
|
||||||
/*
|
/*
|
||||||
|
@ -183,10 +183,15 @@ private void printConfiguredHosts() {
|
|||||||
YarnConfiguration.DEFAULT_RM_NODES_INCLUDE_FILE_PATH) + " out=" +
|
YarnConfiguration.DEFAULT_RM_NODES_INCLUDE_FILE_PATH) + " out=" +
|
||||||
conf.get(YarnConfiguration.RM_NODES_EXCLUDE_FILE_PATH,
|
conf.get(YarnConfiguration.RM_NODES_EXCLUDE_FILE_PATH,
|
||||||
YarnConfiguration.DEFAULT_RM_NODES_EXCLUDE_FILE_PATH));
|
YarnConfiguration.DEFAULT_RM_NODES_EXCLUDE_FILE_PATH));
|
||||||
for (String include : hostsReader.getHosts()) {
|
|
||||||
|
Set<String> hostsList = new HashSet<String>();
|
||||||
|
Set<String> excludeList = new HashSet<String>();
|
||||||
|
hostsReader.getHostDetails(hostsList, excludeList);
|
||||||
|
|
||||||
|
for (String include : hostsList) {
|
||||||
LOG.debug("include: " + include);
|
LOG.debug("include: " + include);
|
||||||
}
|
}
|
||||||
for (String exclude : hostsReader.getExcludedHosts()) {
|
for (String exclude : excludeList) {
|
||||||
LOG.debug("exclude: " + exclude);
|
LOG.debug("exclude: " + exclude);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@ -208,25 +213,17 @@ public void refreshNodes(Configuration yarnConf) throws IOException,
|
|||||||
|
|
||||||
private void refreshHostsReader(Configuration yarnConf) throws IOException,
|
private void refreshHostsReader(Configuration yarnConf) throws IOException,
|
||||||
YarnException {
|
YarnException {
|
||||||
synchronized (hostsReader) {
|
if (null == yarnConf) {
|
||||||
if (null == yarnConf) {
|
yarnConf = new YarnConfiguration();
|
||||||
yarnConf = new YarnConfiguration();
|
|
||||||
}
|
|
||||||
includesFile =
|
|
||||||
yarnConf.get(YarnConfiguration.RM_NODES_INCLUDE_FILE_PATH,
|
|
||||||
YarnConfiguration.DEFAULT_RM_NODES_INCLUDE_FILE_PATH);
|
|
||||||
excludesFile =
|
|
||||||
yarnConf.get(YarnConfiguration.RM_NODES_EXCLUDE_FILE_PATH,
|
|
||||||
YarnConfiguration.DEFAULT_RM_NODES_EXCLUDE_FILE_PATH);
|
|
||||||
hostsReader.updateFileNames(includesFile, excludesFile);
|
|
||||||
hostsReader.refresh(
|
|
||||||
includesFile.isEmpty() ? null : this.rmContext
|
|
||||||
.getConfigurationProvider().getConfigurationInputStream(
|
|
||||||
this.conf, includesFile), excludesFile.isEmpty() ? null
|
|
||||||
: this.rmContext.getConfigurationProvider()
|
|
||||||
.getConfigurationInputStream(this.conf, excludesFile));
|
|
||||||
printConfiguredHosts();
|
|
||||||
}
|
}
|
||||||
|
includesFile =
|
||||||
|
yarnConf.get(YarnConfiguration.RM_NODES_INCLUDE_FILE_PATH,
|
||||||
|
YarnConfiguration.DEFAULT_RM_NODES_INCLUDE_FILE_PATH);
|
||||||
|
excludesFile =
|
||||||
|
yarnConf.get(YarnConfiguration.RM_NODES_EXCLUDE_FILE_PATH,
|
||||||
|
YarnConfiguration.DEFAULT_RM_NODES_EXCLUDE_FILE_PATH);
|
||||||
|
hostsReader.refresh(includesFile, excludesFile);
|
||||||
|
printConfiguredHosts();
|
||||||
}
|
}
|
||||||
|
|
||||||
private void setDecomissionedNMs() {
|
private void setDecomissionedNMs() {
|
||||||
@ -364,13 +361,13 @@ public void run() {
|
|||||||
|
|
||||||
public boolean isValidNode(String hostName) {
|
public boolean isValidNode(String hostName) {
|
||||||
String ip = resolver.resolve(hostName);
|
String ip = resolver.resolve(hostName);
|
||||||
synchronized (hostsReader) {
|
Set<String> hostsList = new HashSet<String>();
|
||||||
Set<String> hostsList = hostsReader.getHosts();
|
Set<String> excludeList = new HashSet<String>();
|
||||||
Set<String> excludeList = hostsReader.getExcludedHosts();
|
hostsReader.getHostDetails(hostsList, excludeList);
|
||||||
return (hostsList.isEmpty() || hostsList.contains(hostName) || hostsList
|
|
||||||
.contains(ip))
|
return (hostsList.isEmpty() || hostsList.contains(hostName) || hostsList
|
||||||
&& !(excludeList.contains(hostName) || excludeList.contains(ip));
|
.contains(ip))
|
||||||
}
|
&& !(excludeList.contains(hostName) || excludeList.contains(ip));
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
@ -467,17 +464,15 @@ private void updateInactiveNodes() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
public boolean isUntrackedNode(String hostName) {
|
public boolean isUntrackedNode(String hostName) {
|
||||||
boolean untracked;
|
|
||||||
String ip = resolver.resolve(hostName);
|
String ip = resolver.resolve(hostName);
|
||||||
|
|
||||||
synchronized (hostsReader) {
|
Set<String> hostsList = new HashSet<String>();
|
||||||
Set<String> hostsList = hostsReader.getHosts();
|
Set<String> excludeList = new HashSet<String>();
|
||||||
Set<String> excludeList = hostsReader.getExcludedHosts();
|
hostsReader.getHostDetails(hostsList, excludeList);
|
||||||
untracked = !hostsList.isEmpty() &&
|
|
||||||
!hostsList.contains(hostName) && !hostsList.contains(ip) &&
|
return !hostsList.isEmpty() && !hostsList.contains(hostName)
|
||||||
!excludeList.contains(hostName) && !excludeList.contains(ip);
|
&& !hostsList.contains(ip) && !excludeList.contains(hostName)
|
||||||
}
|
&& !excludeList.contains(ip);
|
||||||
return untracked;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
Loading…
Reference in New Issue
Block a user