Author: todd
Date: Thu Aug 4 21:56:17 2011
New Revision: 1154029
URL: http://svn.apache.org/viewvc?rev=1154029&view=rev
Log:
HDFS-2225. Refactor file management so it's not in classes which should be
generic. Contributed by Ivan Kelly.
Modified:
hadoop/common/trunk/hdfs/CHANGES.txt
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/BackupJournalManager.java
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FSEditLog.java
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FSImageStorageInspector.java
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FSImageTransactionalStorageInspector.java
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FileJournalManager.java
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/JournalManager.java
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/NNStorageRetentionManager.java
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/FSImageTestUtil.java
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestBackupNode.java
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestCheckPointForSecurityTokens.java
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestFSImageStorageInspector.java
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestNNStorageRetentionManager.java
Modified: hadoop/common/trunk/hdfs/CHANGES.txt
URL:
http://svn.apache.org/viewvc/hadoop/common/trunk/hdfs/CHANGES.txt?rev=1154029&r1=1154028&r2=1154029&view=diff
==============================================================================
--- hadoop/common/trunk/hdfs/CHANGES.txt (original)
+++ hadoop/common/trunk/hdfs/CHANGES.txt Thu Aug 4 21:56:17 2011
@@ -632,6 +632,9 @@ Trunk (unreleased changes)
HDFS-2187. Make EditLogInputStream act like an iterator over FSEditLogOps
(Ivan Kelly and todd via todd)
+ HDFS-2225. Refactor file management so it's not in classes which should
+ be generic. (Ivan Kelly via todd)
+
OPTIMIZATIONS
HDFS-1458. Improve checkpoint performance by avoiding unnecessary image
Modified:
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/BackupJournalManager.java
URL:
http://svn.apache.org/viewvc/hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/BackupJournalManager.java?rev=1154029&r1=1154028&r2=1154029&view=diff
==============================================================================
---
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/BackupJournalManager.java
(original)
+++
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/BackupJournalManager.java
Thu Aug 4 21:56:17 2011
@@ -54,7 +54,7 @@ class BackupJournalManager implements Jo
}
@Override
- public void purgeLogsOlderThan(long minTxIdToKeep, StoragePurger purger)
+ public void purgeLogsOlderThan(long minTxIdToKeep)
throws IOException {
}
Modified:
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FSEditLog.java
URL:
http://svn.apache.org/viewvc/hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FSEditLog.java?rev=1154029&r1=1154028&r2=1154029&view=diff
==============================================================================
---
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FSEditLog.java
(original)
+++
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FSEditLog.java
Thu Aug 4 21:56:17 2011
@@ -862,8 +862,7 @@ public class FSEditLog {
/**
* Archive any log files that are older than the given txid.
*/
- public void purgeLogsOlderThan(
- final long minTxIdToKeep, final StoragePurger purger) {
+ public void purgeLogsOlderThan(final long minTxIdToKeep) {
synchronized (this) {
// synchronized to prevent findbugs warning about inconsistent
// synchronization. This will be JIT-ed out if asserts are
@@ -877,7 +876,7 @@ public class FSEditLog {
mapJournalsAndReportErrors(new JournalClosure() {
@Override
public void apply(JournalAndStream jas) throws IOException {
- jas.manager.purgeLogsOlderThan(minTxIdToKeep, purger);
+ jas.manager.purgeLogsOlderThan(minTxIdToKeep);
}
}, "purging logs older than " + minTxIdToKeep);
}
Modified:
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FSImageStorageInspector.java
URL:
http://svn.apache.org/viewvc/hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FSImageStorageInspector.java?rev=1154029&r1=1154028&r2=1154029&view=diff
==============================================================================
---
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FSImageStorageInspector.java
(original)
+++
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FSImageStorageInspector.java
Thu Aug 4 21:56:17 2011
@@ -96,4 +96,36 @@ abstract class FSImageStorageInspector {
return sb.toString();
}
}
+
+ /**
+ * Record of an image that has been located and had its filename parsed.
+ */
+ static class FSImageFile {
+ final StorageDirectory sd;
+ final long txId;
+ private final File file;
+
+ FSImageFile(StorageDirectory sd, File file, long txId) {
+ assert txId >= 0 : "Invalid txid on " + file +": " + txId;
+
+ this.sd = sd;
+ this.txId = txId;
+ this.file = file;
+ }
+
+ File getFile() {
+ return file;
+ }
+
+ public long getCheckpointTxId() {
+ return txId;
+ }
+
+ @Override
+ public String toString() {
+ return String.format("FSImageFile(file=%s, cpktTxId=%019d)",
+ file.toString(), txId);
+ }
+ }
+
}
Modified:
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FSImageTransactionalStorageInspector.java
URL:
http://svn.apache.org/viewvc/hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FSImageTransactionalStorageInspector.java?rev=1154029&r1=1154028&r2=1154029&view=diff
==============================================================================
---
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FSImageTransactionalStorageInspector.java
(original)
+++
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FSImageTransactionalStorageInspector.java
Thu Aug 4 21:56:17 2011
@@ -37,9 +37,9 @@ import org.apache.commons.logging.LogFac
import org.apache.hadoop.fs.FileUtil;
import org.apache.hadoop.hdfs.protocol.FSConstants;
import org.apache.hadoop.hdfs.server.common.Storage.StorageDirectory;
-import
org.apache.hadoop.hdfs.server.namenode.FSEditLogLoader.EditLogValidation;
import org.apache.hadoop.hdfs.server.namenode.NNStorage.NameNodeDirType;
import org.apache.hadoop.hdfs.server.namenode.NNStorage.NameNodeFile;
+import org.apache.hadoop.hdfs.server.namenode.FileJournalManager.EditLogFile;
import org.apache.hadoop.hdfs.server.protocol.RemoteEditLog;
import org.apache.hadoop.hdfs.server.protocol.RemoteEditLogManifest;
@@ -54,17 +54,13 @@ class FSImageTransactionalStorageInspect
private boolean needToSave = false;
private boolean isUpgradeFinalized = true;
- List<FoundFSImage> foundImages = new ArrayList<FoundFSImage>();
- List<FoundEditLog> foundEditLogs = new ArrayList<FoundEditLog>();
+ List<FSImageFile> foundImages = new ArrayList<FSImageFile>();
+ List<EditLogFile> foundEditLogs = new ArrayList<EditLogFile>();
SortedMap<Long, LogGroup> logGroups = new TreeMap<Long, LogGroup>();
long maxSeenTxId = 0;
private static final Pattern IMAGE_REGEX = Pattern.compile(
NameNodeFile.IMAGE.getName() + "_(\\d+)");
- private static final Pattern EDITS_REGEX = Pattern.compile(
- NameNodeFile.EDITS.getName() + "_(\\d+)-(\\d+)");
- private static final Pattern EDITS_INPROGRESS_REGEX = Pattern.compile(
- NameNodeFile.EDITS_INPROGRESS.getName() + "_(\\d+)");
@Override
public void inspectDirectory(StorageDirectory sd) throws IOException {
@@ -95,7 +91,7 @@ class FSImageTransactionalStorageInspect
if (sd.getStorageDirType().isOfType(NameNodeDirType.IMAGE)) {
try {
long txid = Long.valueOf(imageMatch.group(1));
- foundImages.add(new FoundFSImage(sd, f, txid));
+ foundImages.add(new FSImageFile(sd, f, txid));
} catch (NumberFormatException nfe) {
LOG.error("Image file " + f + " has improperly formatted " +
"transaction ID");
@@ -117,9 +113,10 @@ class FSImageTransactionalStorageInspect
LOG.warn("Unable to determine the max transaction ID seen by " + sd,
ioe);
}
- List<FoundEditLog> editLogs = matchEditLogs(filesInStorage);
+ List<EditLogFile> editLogs
+ = FileJournalManager.matchEditLogs(filesInStorage);
if (sd.getStorageDirType().isOfType(NameNodeDirType.EDITS)) {
- for (FoundEditLog log : editLogs) {
+ for (EditLogFile log : editLogs) {
addEditLog(log);
}
} else if (!editLogs.isEmpty()){
@@ -133,47 +130,12 @@ class FSImageTransactionalStorageInspect
isUpgradeFinalized = isUpgradeFinalized && !sd.getPreviousDir().exists();
}
- static List<FoundEditLog> matchEditLogs(File[] filesInStorage) {
- List<FoundEditLog> ret = Lists.newArrayList();
- for (File f : filesInStorage) {
- String name = f.getName();
- // Check for edits
- Matcher editsMatch = EDITS_REGEX.matcher(name);
- if (editsMatch.matches()) {
- try {
- long startTxId = Long.valueOf(editsMatch.group(1));
- long endTxId = Long.valueOf(editsMatch.group(2));
- ret.add(new FoundEditLog(f, startTxId, endTxId));
- } catch (NumberFormatException nfe) {
- LOG.error("Edits file " + f + " has improperly formatted " +
- "transaction ID");
- // skip
- }
- }
-
- // Check for in-progress edits
- Matcher inProgressEditsMatch = EDITS_INPROGRESS_REGEX.matcher(name);
- if (inProgressEditsMatch.matches()) {
- try {
- long startTxId = Long.valueOf(inProgressEditsMatch.group(1));
- ret.add(
- new FoundEditLog(f, startTxId, FoundEditLog.UNKNOWN_END));
- } catch (NumberFormatException nfe) {
- LOG.error("In-progress edits file " + f + " has improperly " +
- "formatted transaction ID");
- // skip
- }
- }
- }
- return ret;
- }
-
- private void addEditLog(FoundEditLog foundEditLog) {
+ private void addEditLog(EditLogFile foundEditLog) {
foundEditLogs.add(foundEditLog);
- LogGroup group = logGroups.get(foundEditLog.startTxId);
+ LogGroup group = logGroups.get(foundEditLog.getFirstTxId());
if (group == null) {
- group = new LogGroup(foundEditLog.startTxId);
- logGroups.put(foundEditLog.startTxId, group);
+ group = new LogGroup(foundEditLog.getFirstTxId());
+ logGroups.put(foundEditLog.getFirstTxId(), group);
}
group.add(foundEditLog);
}
@@ -191,9 +153,9 @@ class FSImageTransactionalStorageInspect
*
* Returns null if no images were found.
*/
- FoundFSImage getLatestImage() {
- FoundFSImage ret = null;
- for (FoundFSImage img : foundImages) {
+ FSImageFile getLatestImage() {
+ FSImageFile ret = null;
+ for (FSImageFile img : foundImages) {
if (ret == null || img.txId > ret.txId) {
ret = img;
}
@@ -201,11 +163,11 @@ class FSImageTransactionalStorageInspect
return ret;
}
- public List<FoundFSImage> getFoundImages() {
+ public List<FSImageFile> getFoundImages() {
return ImmutableList.copyOf(foundImages);
}
- public List<FoundEditLog> getFoundEditLogs() {
+ public List<EditLogFile> getEditLogFiles() {
return ImmutableList.copyOf(foundEditLogs);
}
@@ -215,7 +177,7 @@ class FSImageTransactionalStorageInspect
throw new FileNotFoundException("No valid image files found");
}
- FoundFSImage recoveryImage = getLatestImage();
+ FSImageFile recoveryImage = getLatestImage();
LogLoadPlan logPlan = createLogLoadPlan(recoveryImage.txId,
Long.MAX_VALUE);
return new TransactionalLoadPlan(recoveryImage,
@@ -233,7 +195,7 @@ class FSImageTransactionalStorageInspect
LogLoadPlan createLogLoadPlan(long sinceTxId, long maxStartTxId) throws
IOException {
long expectedTxId = sinceTxId + 1;
- List<FoundEditLog> recoveryLogs = new ArrayList<FoundEditLog>();
+ List<EditLogFile> recoveryLogs = new ArrayList<EditLogFile>();
SortedMap<Long, LogGroup> tailGroups = logGroups.tailMap(expectedTxId);
if (logGroups.size() > tailGroups.size()) {
@@ -312,10 +274,10 @@ class FSImageTransactionalStorageInspect
for (LogGroup g : logGroups.values()) {
if (!g.hasFinalized) continue;
- FoundEditLog fel = g.getBestNonCorruptLog();
+ EditLogFile fel = g.getBestNonCorruptLog();
if (fel.getLastTxId() < sinceTxId) continue;
- logs.add(new RemoteEditLog(fel.getStartTxId(),
+ logs.add(new RemoteEditLog(fel.getFirstTxId(),
fel.getLastTxId()));
}
@@ -330,7 +292,7 @@ class FSImageTransactionalStorageInspect
*/
static class LogGroup {
long startTxId;
- List<FoundEditLog> logs = new ArrayList<FoundEditLog>();;
+ List<EditLogFile> logs = new ArrayList<EditLogFile>();;
private Set<Long> endTxIds = new TreeSet<Long>();
private boolean hasInProgress = false;
private boolean hasFinalized = false;
@@ -339,15 +301,15 @@ class FSImageTransactionalStorageInspect
this.startTxId = startTxId;
}
- FoundEditLog getBestNonCorruptLog() {
+ EditLogFile getBestNonCorruptLog() {
// First look for non-corrupt finalized logs
- for (FoundEditLog log : logs) {
+ for (EditLogFile log : logs) {
if (!log.isCorrupt() && !log.isInProgress()) {
return log;
}
}
// Then look for non-corrupt in-progress logs
- for (FoundEditLog log : logs) {
+ for (EditLogFile log : logs) {
if (!log.isCorrupt()) {
return log;
}
@@ -364,7 +326,7 @@ class FSImageTransactionalStorageInspect
* @return true if we can determine the last txid in this log group.
*/
boolean hasKnownLastTxId() {
- for (FoundEditLog log : logs) {
+ for (EditLogFile log : logs) {
if (!log.isInProgress()) {
return true;
}
@@ -378,24 +340,24 @@ class FSImageTransactionalStorageInspect
* {@see #hasKnownLastTxId()}
*/
long getLastTxId() {
- for (FoundEditLog log : logs) {
+ for (EditLogFile log : logs) {
if (!log.isInProgress()) {
- return log.lastTxId;
+ return log.getLastTxId();
}
}
throw new IllegalStateException("LogGroup only has in-progress logs");
}
- void add(FoundEditLog log) {
- assert log.getStartTxId() == startTxId;
+ void add(EditLogFile log) {
+ assert log.getFirstTxId() == startTxId;
logs.add(log);
if (log.isInProgress()) {
hasInProgress = true;
} else {
hasFinalized = true;
- endTxIds.add(log.lastTxId);
+ endTxIds.add(log.getLastTxId());
}
}
@@ -422,7 +384,7 @@ class FSImageTransactionalStorageInspect
* The in-progress logs in this case should be considered corrupt.
*/
private void planMixedLogRecovery() throws IOException {
- for (FoundEditLog log : logs) {
+ for (EditLogFile log : logs) {
if (log.isInProgress()) {
LOG.warn("Log at " + log.getFile() + " is in progress, but " +
"other logs starting at the same txid " + startTxId +
@@ -446,7 +408,7 @@ class FSImageTransactionalStorageInspect
"crash)");
if (logs.size() == 1) {
// Only one log, it's our only choice!
- FoundEditLog log = logs.get(0);
+ EditLogFile log = logs.get(0);
if (log.validateLog().numTransactions == 0) {
// If it has no transactions, we should consider it corrupt just
// to be conservative.
@@ -459,7 +421,7 @@ class FSImageTransactionalStorageInspect
}
long maxValidTxnCount = Long.MIN_VALUE;
- for (FoundEditLog log : logs) {
+ for (EditLogFile log : logs) {
long validTxnCount = log.validateLog().numTransactions;
LOG.warn(" Log " + log.getFile() +
" valid txns=" + validTxnCount +
@@ -467,7 +429,7 @@ class FSImageTransactionalStorageInspect
maxValidTxnCount = Math.max(maxValidTxnCount, validTxnCount);
}
- for (FoundEditLog log : logs) {
+ for (EditLogFile log : logs) {
long txns = log.validateLog().numTransactions;
if (txns < maxValidTxnCount) {
LOG.warn("Marking log at " + log.getFile() + " as corrupt since " +
@@ -499,7 +461,7 @@ class FSImageTransactionalStorageInspect
}
void recover() throws IOException {
- for (FoundEditLog log : logs) {
+ for (EditLogFile log : logs) {
if (log.isCorrupt()) {
log.moveAsideCorruptFile();
} else if (log.isInProgress()) {
@@ -508,131 +470,12 @@ class FSImageTransactionalStorageInspect
}
}
}
-
- /**
- * Record of an image that has been located and had its filename parsed.
- */
- static class FoundFSImage {
- final StorageDirectory sd;
- final long txId;
- private final File file;
-
- FoundFSImage(StorageDirectory sd, File file, long txId) {
- assert txId >= 0 : "Invalid txid on " + file +": " + txId;
-
- this.sd = sd;
- this.txId = txId;
- this.file = file;
- }
-
- File getFile() {
- return file;
- }
-
- public long getTxId() {
- return txId;
- }
-
- @Override
- public String toString() {
- return file.toString();
- }
- }
- /**
- * Record of an edit log that has been located and had its filename parsed.
- */
- static class FoundEditLog {
- File file;
- final long startTxId;
- long lastTxId;
-
- private EditLogValidation cachedValidation = null;
- private boolean isCorrupt = false;
-
- static final long UNKNOWN_END = -1;
-
- FoundEditLog(File file,
- long startTxId, long endTxId) {
- assert endTxId == UNKNOWN_END || endTxId >= startTxId;
- assert startTxId > 0;
- assert file != null;
-
- this.startTxId = startTxId;
- this.lastTxId = endTxId;
- this.file = file;
- }
-
- public void finalizeLog() throws IOException {
- long numTransactions = validateLog().numTransactions;
- long lastTxId = startTxId + numTransactions - 1;
- File dst = new File(file.getParentFile(),
- NNStorage.getFinalizedEditsFileName(startTxId, lastTxId));
- LOG.info("Finalizing edits log " + file + " by renaming to "
- + dst.getName());
- if (!file.renameTo(dst)) {
- throw new IOException("Couldn't finalize log " +
- file + " to " + dst);
- }
- this.lastTxId = lastTxId;
- file = dst;
- }
-
- long getStartTxId() {
- return startTxId;
- }
-
- long getLastTxId() {
- return lastTxId;
- }
-
- EditLogValidation validateLog() throws IOException {
- if (cachedValidation == null) {
- cachedValidation = EditLogFileInputStream.validateEditLog(file);
- }
- return cachedValidation;
- }
-
- boolean isInProgress() {
- return (lastTxId == UNKNOWN_END);
- }
-
- File getFile() {
- return file;
- }
-
- void markCorrupt() {
- isCorrupt = true;
- }
-
- boolean isCorrupt() {
- return isCorrupt;
- }
-
- void moveAsideCorruptFile() throws IOException {
- assert isCorrupt;
-
- File src = file;
- File dst = new File(src.getParent(), src.getName() + ".corrupt");
- boolean success = src.renameTo(dst);
- if (!success) {
- throw new IOException(
- "Couldn't rename corrupt log " + src + " to " + dst);
- }
- file = dst;
- }
-
- @Override
- public String toString() {
- return file.toString();
- }
- }
-
static class TransactionalLoadPlan extends LoadPlan {
- final FoundFSImage image;
+ final FSImageFile image;
final LogLoadPlan logPlan;
- public TransactionalLoadPlan(FoundFSImage image,
+ public TransactionalLoadPlan(FSImageFile image,
LogLoadPlan logPlan) {
super();
this.image = image;
@@ -662,10 +505,10 @@ class FSImageTransactionalStorageInspect
}
static class LogLoadPlan {
- final List<FoundEditLog> editLogs;
+ final List<EditLogFile> editLogs;
final List<LogGroup> logGroupsToRecover;
- LogLoadPlan(List<FoundEditLog> editLogs,
+ LogLoadPlan(List<EditLogFile> editLogs,
List<LogGroup> logGroupsToRecover) {
this.editLogs = editLogs;
this.logGroupsToRecover = logGroupsToRecover;
@@ -679,7 +522,7 @@ class FSImageTransactionalStorageInspect
public List<File> getEditsFiles() {
List<File> ret = new ArrayList<File>();
- for (FoundEditLog log : editLogs) {
+ for (EditLogFile log : editLogs) {
ret.add(log.getFile());
}
return ret;
Modified:
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FileJournalManager.java
URL:
http://svn.apache.org/viewvc/hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FileJournalManager.java?rev=1154029&r1=1154028&r2=1154029&view=diff
==============================================================================
---
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FileJournalManager.java
(original)
+++
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/FileJournalManager.java
Thu Aug 4 21:56:17 2011
@@ -23,14 +23,20 @@ import org.apache.commons.logging.LogFac
import java.io.File;
import java.io.IOException;
import java.util.List;
+import java.util.Comparator;
+import java.util.regex.Matcher;
+import java.util.regex.Pattern;
import org.apache.hadoop.fs.FileUtil;
import org.apache.hadoop.hdfs.server.common.Storage.StorageDirectory;
-import
org.apache.hadoop.hdfs.server.namenode.FSImageTransactionalStorageInspector.FoundEditLog;
import
org.apache.hadoop.hdfs.server.namenode.NNStorageRetentionManager.StoragePurger;
+import
org.apache.hadoop.hdfs.server.namenode.FSEditLogLoader.EditLogValidation;
+import org.apache.hadoop.hdfs.server.namenode.NNStorage.NameNodeFile;
import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Preconditions;
+import com.google.common.collect.Lists;
+import com.google.common.collect.ComparisonChain;
/**
* Journal manager for the common case of edits files being written
@@ -45,6 +51,15 @@ class FileJournalManager implements Jour
private final StorageDirectory sd;
private int outputBufferCapacity = 512*1024;
+ private static final Pattern EDITS_REGEX = Pattern.compile(
+ NameNodeFile.EDITS.getName() + "_(\\d+)-(\\d+)");
+ private static final Pattern EDITS_INPROGRESS_REGEX = Pattern.compile(
+ NameNodeFile.EDITS_INPROGRESS.getName() + "_(\\d+)");
+
+ @VisibleForTesting
+ StoragePurger purger
+ = new NNStorageRetentionManager.DeletionStoragePurger();
+
public FileJournalManager(StorageDirectory sd) {
this.sd = sd;
}
@@ -91,13 +106,13 @@ class FileJournalManager implements Jour
}
@Override
- public void purgeLogsOlderThan(long minTxIdToKeep, StoragePurger purger)
+ public void purgeLogsOlderThan(long minTxIdToKeep)
throws IOException {
File[] files = FileUtil.listFiles(sd.getCurrentDir());
- List<FoundEditLog> editLogs =
- FSImageTransactionalStorageInspector.matchEditLogs(files);
- for (FoundEditLog log : editLogs) {
- if (log.getStartTxId() < minTxIdToKeep &&
+ List<EditLogFile> editLogs =
+ FileJournalManager.matchEditLogs(files);
+ for (EditLogFile log : editLogs) {
+ if (log.getFirstTxId() < minTxIdToKeep &&
log.getLastTxId() < minTxIdToKeep) {
purger.purgeLog(log);
}
@@ -111,4 +126,139 @@ class FileJournalManager implements Jour
return new EditLogFileInputStream(f);
}
+ static List<EditLogFile> matchEditLogs(File[] filesInStorage) {
+ List<EditLogFile> ret = Lists.newArrayList();
+ for (File f : filesInStorage) {
+ String name = f.getName();
+ // Check for edits
+ Matcher editsMatch = EDITS_REGEX.matcher(name);
+ if (editsMatch.matches()) {
+ try {
+ long startTxId = Long.valueOf(editsMatch.group(1));
+ long endTxId = Long.valueOf(editsMatch.group(2));
+ ret.add(new EditLogFile(f, startTxId, endTxId));
+ } catch (NumberFormatException nfe) {
+ LOG.error("Edits file " + f + " has improperly formatted " +
+ "transaction ID");
+ // skip
+ }
+ }
+
+ // Check for in-progress edits
+ Matcher inProgressEditsMatch = EDITS_INPROGRESS_REGEX.matcher(name);
+ if (inProgressEditsMatch.matches()) {
+ try {
+ long startTxId = Long.valueOf(inProgressEditsMatch.group(1));
+ ret.add(
+ new EditLogFile(f, startTxId, EditLogFile.UNKNOWN_END));
+ } catch (NumberFormatException nfe) {
+ LOG.error("In-progress edits file " + f + " has improperly " +
+ "formatted transaction ID");
+ // skip
+ }
+ }
+ }
+ return ret;
+ }
+
+ /**
+ * Record of an edit log that has been located and had its filename parsed.
+ */
+ static class EditLogFile {
+ private File file;
+ private final long firstTxId;
+ private long lastTxId;
+
+ private EditLogValidation cachedValidation = null;
+ private boolean isCorrupt = false;
+
+ static final long UNKNOWN_END = -1;
+
+ final static Comparator<EditLogFile> COMPARE_BY_START_TXID
+ = new Comparator<EditLogFile>() {
+ public int compare(EditLogFile a, EditLogFile b) {
+ return ComparisonChain.start()
+ .compare(a.getFirstTxId(), b.getFirstTxId())
+ .compare(a.getLastTxId(), b.getLastTxId())
+ .result();
+ }
+ };
+
+ EditLogFile(File file,
+ long firstTxId, long lastTxId) {
+ assert lastTxId == UNKNOWN_END || lastTxId >= firstTxId;
+ assert firstTxId > 0;
+ assert file != null;
+
+ this.firstTxId = firstTxId;
+ this.lastTxId = lastTxId;
+ this.file = file;
+ }
+
+ public void finalizeLog() throws IOException {
+ long numTransactions = validateLog().numTransactions;
+ long lastTxId = firstTxId + numTransactions - 1;
+ File dst = new File(file.getParentFile(),
+ NNStorage.getFinalizedEditsFileName(firstTxId, lastTxId));
+ LOG.info("Finalizing edits log " + file + " by renaming to "
+ + dst.getName());
+ if (!file.renameTo(dst)) {
+ throw new IOException("Couldn't finalize log " +
+ file + " to " + dst);
+ }
+ this.lastTxId = lastTxId;
+ file = dst;
+ }
+
+ long getFirstTxId() {
+ return firstTxId;
+ }
+
+ long getLastTxId() {
+ return lastTxId;
+ }
+
+ EditLogValidation validateLog() throws IOException {
+ if (cachedValidation == null) {
+ cachedValidation = EditLogFileInputStream.validateEditLog(file);
+ }
+ return cachedValidation;
+ }
+
+ boolean isInProgress() {
+ return (lastTxId == UNKNOWN_END);
+ }
+
+ File getFile() {
+ return file;
+ }
+
+ void markCorrupt() {
+ isCorrupt = true;
+ }
+
+ boolean isCorrupt() {
+ return isCorrupt;
+ }
+
+ void moveAsideCorruptFile() throws IOException {
+ assert isCorrupt;
+
+ File src = file;
+ File dst = new File(src.getParent(), src.getName() + ".corrupt");
+ boolean success = src.renameTo(dst);
+ if (!success) {
+ throw new IOException(
+ "Couldn't rename corrupt log " + src + " to " + dst);
+ }
+ file = dst;
+ }
+
+ @Override
+ public String toString() {
+ return String.format("EditLogFile(file=%s,first=%019d,last=%019d,"
+ +"inProgress=%b,corrupt=%b)", file.toString(),
+ firstTxId, lastTxId, isInProgress(), isCorrupt);
+ }
+ }
}
Modified:
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/JournalManager.java
URL:
http://svn.apache.org/viewvc/hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/JournalManager.java?rev=1154029&r1=1154028&r2=1154029&view=diff
==============================================================================
---
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/JournalManager.java
(original)
+++
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/JournalManager.java
Thu Aug 4 21:56:17 2011
@@ -55,7 +55,7 @@ interface JournalManager {
* @param purger the purging implementation to use
* @throws IOException if purging fails
*/
- void purgeLogsOlderThan(long minTxIdToKeep, StoragePurger purger)
+ void purgeLogsOlderThan(long minTxIdToKeep)
throws IOException;
/**
Modified:
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/NNStorageRetentionManager.java
URL:
http://svn.apache.org/viewvc/hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/NNStorageRetentionManager.java?rev=1154029&r1=1154028&r2=1154029&view=diff
==============================================================================
---
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/NNStorageRetentionManager.java
(original)
+++
hadoop/common/trunk/hdfs/src/java/org/apache/hadoop/hdfs/server/namenode/NNStorageRetentionManager.java
Thu Aug 4 21:56:17 2011
@@ -27,8 +27,8 @@ import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.hdfs.DFSConfigKeys;
-import
org.apache.hadoop.hdfs.server.namenode.FSImageTransactionalStorageInspector.FoundEditLog;
-import
org.apache.hadoop.hdfs.server.namenode.FSImageTransactionalStorageInspector.FoundFSImage;
+import
org.apache.hadoop.hdfs.server.namenode.FSImageStorageInspector.FSImageFile;
+import org.apache.hadoop.hdfs.server.namenode.FileJournalManager.EditLogFile;
import org.apache.hadoop.hdfs.util.MD5FileUtils;
import com.google.common.collect.Lists;
@@ -80,14 +80,14 @@ public class NNStorageRetentionManager {
// If fsimage_N is the image we want to keep, then we need to keep
// all txns > N. We can remove anything < N+1, since fsimage_N
// reflects the state up to and including N.
- editLog.purgeLogsOlderThan(minImageTxId + 1, purger);
+ editLog.purgeLogsOlderThan(minImageTxId + 1);
}
private void purgeCheckpointsOlderThan(
FSImageTransactionalStorageInspector inspector,
long minTxId) {
- for (FoundFSImage image : inspector.getFoundImages()) {
- if (image.getTxId() < minTxId) {
+ for (FSImageFile image : inspector.getFoundImages()) {
+ if (image.getCheckpointTxId() < minTxId) {
LOG.info("Purging old image " + image);
purger.purgeImage(image);
}
@@ -101,10 +101,10 @@ public class NNStorageRetentionManager {
*/
private long getImageTxIdToRetain(FSImageTransactionalStorageInspector
inspector) {
- List<FoundFSImage> images = inspector.getFoundImages();
+ List<FSImageFile> images = inspector.getFoundImages();
TreeSet<Long> imageTxIds = Sets.newTreeSet();
- for (FoundFSImage image : images) {
- imageTxIds.add(image.getTxId());
+ for (FSImageFile image : images) {
+ imageTxIds.add(image.getCheckpointTxId());
}
List<Long> imageTxIdsList = Lists.newArrayList(imageTxIds);
@@ -124,18 +124,18 @@ public class NNStorageRetentionManager {
* Interface responsible for disposing of old checkpoints and edit logs.
*/
static interface StoragePurger {
- void purgeLog(FoundEditLog log);
- void purgeImage(FoundFSImage image);
+ void purgeLog(EditLogFile log);
+ void purgeImage(FSImageFile image);
}
static class DeletionStoragePurger implements StoragePurger {
@Override
- public void purgeLog(FoundEditLog log) {
+ public void purgeLog(EditLogFile log) {
deleteOrWarn(log.getFile());
}
@Override
- public void purgeImage(FoundFSImage image) {
+ public void purgeImage(FSImageFile image) {
deleteOrWarn(image.getFile());
deleteOrWarn(MD5FileUtils.getDigestFileForFile(image.getFile()));
}
Modified:
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/FSImageTestUtil.java
URL:
http://svn.apache.org/viewvc/hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/FSImageTestUtil.java?rev=1154029&r1=1154028&r2=1154029&view=diff
==============================================================================
---
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/FSImageTestUtil.java
(original)
+++
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/FSImageTestUtil.java
Thu Aug 4 21:56:17 2011
@@ -25,7 +25,6 @@ import java.io.RandomAccessFile;
import java.net.URI;
import java.util.ArrayList;
import java.util.Collections;
-import java.util.Comparator;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -34,15 +33,14 @@ import java.util.Set;
import org.apache.hadoop.hdfs.MiniDFSCluster;
import org.apache.hadoop.hdfs.server.common.Storage.StorageDirectory;
-import
org.apache.hadoop.hdfs.server.namenode.FSImageTransactionalStorageInspector.FoundEditLog;
-import
org.apache.hadoop.hdfs.server.namenode.FSImageTransactionalStorageInspector.FoundFSImage;
+import org.apache.hadoop.hdfs.server.namenode.FileJournalManager.EditLogFile;
+import
org.apache.hadoop.hdfs.server.namenode.FSImageStorageInspector.FSImageFile;
import org.apache.hadoop.hdfs.server.namenode.NNStorage.NameNodeDirType;
import org.apache.hadoop.hdfs.util.MD5FileUtils;
import org.apache.hadoop.io.IOUtils;
import org.mockito.Mockito;
import com.google.common.base.Joiner;
-import com.google.common.collect.ComparisonChain;
import com.google.common.collect.Lists;
import com.google.common.collect.Maps;
import com.google.common.collect.Sets;
@@ -154,9 +152,9 @@ public abstract class FSImageTestUtil {
for (File dir : dirs) {
FSImageTransactionalStorageInspector inspector =
inspectStorageDirectory(dir, NameNodeDirType.IMAGE);
- FoundFSImage latestImage = inspector.getLatestImage();
+ FSImageFile latestImage = inspector.getLatestImage();
assertNotNull("No image in " + dir, latestImage);
- long thisTxId = latestImage.getTxId();
+ long thisTxId = latestImage.getCheckpointTxId();
if (imageTxId != -1 && thisTxId != imageTxId) {
fail("Storage directory " + dir + " does not have the same " +
"last image index " + imageTxId + " as another");
@@ -283,7 +281,7 @@ public abstract class FSImageTestUtil {
new FSImageTransactionalStorageInspector();
inspector.inspectDirectory(sd);
- FoundFSImage latestImage = inspector.getLatestImage();
+ FSImageFile latestImage = inspector.getLatestImage();
return (latestImage == null) ? null : latestImage.getFile();
}
@@ -316,23 +314,15 @@ public abstract class FSImageTestUtil {
* @return the latest edits log, finalized or otherwise, from the given
* storage directory.
*/
- public static FoundEditLog findLatestEditsLog(StorageDirectory sd)
+ public static EditLogFile findLatestEditsLog(StorageDirectory sd)
throws IOException {
FSImageTransactionalStorageInspector inspector =
new FSImageTransactionalStorageInspector();
inspector.inspectDirectory(sd);
- List<FoundEditLog> foundEditLogs = Lists.newArrayList(
- inspector.getFoundEditLogs());
- return Collections.max(foundEditLogs, new Comparator<FoundEditLog>() {
- @Override
- public int compare(FoundEditLog a, FoundEditLog b) {
- return ComparisonChain.start()
- .compare(a.getStartTxId(), b.getStartTxId())
- .compare(a.getLastTxId(), b.getLastTxId())
- .result();
- }
- });
+ List<EditLogFile> foundEditLogs = Lists.newArrayList(
+ inspector.getEditLogFiles());
+ return Collections.max(foundEditLogs, EditLogFile.COMPARE_BY_START_TXID);
}
/**
Modified:
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestBackupNode.java
URL:
http://svn.apache.org/viewvc/hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestBackupNode.java?rev=1154029&r1=1154028&r2=1154029&view=diff
==============================================================================
---
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestBackupNode.java
(original)
+++
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestBackupNode.java
Thu Aug 4 21:56:17 2011
@@ -33,7 +33,7 @@ import org.apache.hadoop.hdfs.HdfsConfig
import org.apache.hadoop.hdfs.MiniDFSCluster;
import org.apache.hadoop.hdfs.server.common.HdfsConstants.StartupOption;
import org.apache.hadoop.hdfs.server.common.Storage.StorageDirectory;
-import
org.apache.hadoop.hdfs.server.namenode.FSImageTransactionalStorageInspector.FoundEditLog;
+import org.apache.hadoop.hdfs.server.namenode.FileJournalManager.EditLogFile;
import org.apache.hadoop.hdfs.DFSConfigKeys;
import org.apache.hadoop.test.GenericTestUtils;
import org.apache.log4j.Level;
@@ -163,8 +163,8 @@ public class TestBackupNode extends Test
// When shutting down the BN, it shouldn't finalize logs that are
// still open on the NN
- FoundEditLog editsLog = FSImageTestUtil.findLatestEditsLog(sd);
- assertEquals(editsLog.getStartTxId(),
+ EditLogFile editsLog = FSImageTestUtil.findLatestEditsLog(sd);
+ assertEquals(editsLog.getFirstTxId(),
nn.getFSImage().getEditLog().getCurSegmentTxId());
assertTrue("Should not have finalized " + editsLog,
editsLog.isInProgress());
Modified:
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestCheckPointForSecurityTokens.java
URL:
http://svn.apache.org/viewvc/hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestCheckPointForSecurityTokens.java?rev=1154029&r1=1154028&r2=1154029&view=diff
==============================================================================
---
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestCheckPointForSecurityTokens.java
(original)
+++
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestCheckPointForSecurityTokens.java
Thu Aug 4 21:56:17 2011
@@ -28,7 +28,7 @@ import org.apache.hadoop.hdfs.MiniDFSClu
import org.apache.hadoop.hdfs.protocol.FSConstants.SafeModeAction;
import
org.apache.hadoop.hdfs.security.token.delegation.DelegationTokenIdentifier;
import org.apache.hadoop.hdfs.server.common.Storage.StorageDirectory;
-import
org.apache.hadoop.hdfs.server.namenode.FSImageTransactionalStorageInspector.FoundEditLog;
+import org.apache.hadoop.hdfs.server.namenode.FileJournalManager.EditLogFile;
import org.apache.hadoop.hdfs.tools.DFSAdmin;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.security.UserGroupInformation;
@@ -82,7 +82,7 @@ public class TestCheckPointForSecurityTo
// verify that the edits file is NOT empty
NameNode nn = cluster.getNameNode();
for (StorageDirectory sd :
nn.getFSImage().getStorage().dirIterable(null)) {
- FoundEditLog log = FSImageTestUtil.findLatestEditsLog(sd);
+ EditLogFile log = FSImageTestUtil.findLatestEditsLog(sd);
assertTrue(log.isInProgress());
assertEquals("In-progress log " + log + " should have 5 transactions",
5, log.validateLog().numTransactions);
@@ -97,7 +97,7 @@ public class TestCheckPointForSecurityTo
}
// verify that the edits file is empty except for the START txn
for (StorageDirectory sd :
nn.getFSImage().getStorage().dirIterable(null)) {
- FoundEditLog log = FSImageTestUtil.findLatestEditsLog(sd);
+ EditLogFile log = FSImageTestUtil.findLatestEditsLog(sd);
assertTrue(log.isInProgress());
assertEquals("In-progress log " + log + " should only have START txn",
1, log.validateLog().numTransactions);
Modified:
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestFSImageStorageInspector.java
URL:
http://svn.apache.org/viewvc/hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestFSImageStorageInspector.java?rev=1154029&r1=1154028&r2=1154029&view=diff
==============================================================================
---
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestFSImageStorageInspector.java
(original)
+++
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestFSImageStorageInspector.java
Thu Aug 4 21:56:17 2011
@@ -36,8 +36,8 @@ import static org.apache.hadoop.hdfs.ser
import static
org.apache.hadoop.hdfs.server.namenode.NNStorage.getFinalizedEditsFileName;
import static
org.apache.hadoop.hdfs.server.namenode.NNStorage.getImageFileName;
-import
org.apache.hadoop.hdfs.server.namenode.FSImageTransactionalStorageInspector.FoundEditLog;
-import
org.apache.hadoop.hdfs.server.namenode.FSImageTransactionalStorageInspector.FoundFSImage;
+import org.apache.hadoop.hdfs.server.namenode.FileJournalManager.EditLogFile;
+import
org.apache.hadoop.hdfs.server.namenode.FSImageStorageInspector.FSImageFile;
import
org.apache.hadoop.hdfs.server.namenode.FSImageTransactionalStorageInspector.TransactionalLoadPlan;
import
org.apache.hadoop.hdfs.server.namenode.FSImageTransactionalStorageInspector.LogGroup;
import org.apache.hadoop.hdfs.server.namenode.FSImageStorageInspector.LoadPlan;
@@ -72,7 +72,7 @@ public class TestFSImageStorageInspector
assertEquals(2, inspector.foundImages.size());
assertTrue(inspector.foundEditLogs.get(1).isInProgress());
- FoundFSImage latestImage = inspector.getLatestImage();
+ FSImageFile latestImage = inspector.getLatestImage();
assertEquals(456, latestImage.txId);
assertSame(mockDir, latestImage.sd);
assertTrue(inspector.isUpgradeFinalized());
@@ -203,7 +203,7 @@ public class TestFSImageStorageInspector
LogGroup lg = inspector.logGroups.get(123L);
assertEquals(3, lg.logs.size());
- FoundEditLog inProgressLog = lg.logs.get(2);
+ EditLogFile inProgressLog = lg.logs.get(2);
assertTrue(inProgressLog.isInProgress());
LoadPlan plan = inspector.createLoadPlan();
@@ -282,7 +282,7 @@ public class TestFSImageStorageInspector
assertTrue(lg.logs.get(2).isCorrupt());
// Calling recover should move it aside
- FoundEditLog badLog = lg.logs.get(2);
+ EditLogFile badLog = lg.logs.get(2);
Mockito.doNothing().when(badLog).moveAsideCorruptFile();
Mockito.doNothing().when(lg.logs.get(0)).finalizeLog();
Mockito.doNothing().when(lg.logs.get(1)).finalizeLog();
@@ -303,12 +303,12 @@ public class TestFSImageStorageInspector
String path, int numValidTransactions) throws IOException {
for (LogGroup lg : inspector.logGroups.values()) {
- List<FoundEditLog> logs = lg.logs;
+ List<EditLogFile> logs = lg.logs;
for (int i = 0; i < logs.size(); i++) {
- FoundEditLog log = logs.get(i);
- if (log.file.getPath().equals(path)) {
+ EditLogFile log = logs.get(i);
+ if (log.getFile().getPath().equals(path)) {
// mock out its validation
- FoundEditLog spyLog = spy(log);
+ EditLogFile spyLog = spy(log);
doReturn(new FSEditLogLoader.EditLogValidation(-1,
numValidTransactions))
.when(spyLog).validateLog();
logs.set(i, spyLog);
@@ -356,7 +356,7 @@ public class TestFSImageStorageInspector
// Check plan
TransactionalLoadPlan plan =
(TransactionalLoadPlan)inspector.createLoadPlan();
- FoundFSImage pickedImage = plan.image;
+ FSImageFile pickedImage = plan.image;
assertEquals(456, pickedImage.txId);
assertSame(mockImageDir2, pickedImage.sd);
assertEquals(new File("/foo2/current/" + getImageFileName(456)),
Modified:
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestNNStorageRetentionManager.java
URL:
http://svn.apache.org/viewvc/hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestNNStorageRetentionManager.java?rev=1154029&r1=1154028&r2=1154029&view=diff
==============================================================================
---
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestNNStorageRetentionManager.java
(original)
+++
hadoop/common/trunk/hdfs/src/test/hdfs/org/apache/hadoop/hdfs/server/namenode/TestNNStorageRetentionManager.java
Thu Aug 4 21:56:17 2011
@@ -24,8 +24,8 @@ import java.util.Set;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.hdfs.server.common.Storage.StorageDirectory;
-import
org.apache.hadoop.hdfs.server.namenode.FSImageTransactionalStorageInspector.FoundEditLog;
-import
org.apache.hadoop.hdfs.server.namenode.FSImageTransactionalStorageInspector.FoundFSImage;
+import org.apache.hadoop.hdfs.server.namenode.FileJournalManager.EditLogFile;
+import
org.apache.hadoop.hdfs.server.namenode.FSImageStorageInspector.FSImageFile;
import org.apache.hadoop.hdfs.server.namenode.NNStorage.NameNodeDirType;
import static
org.apache.hadoop.hdfs.server.namenode.NNStorage.getInProgressEditsFileName;
import static
org.apache.hadoop.hdfs.server.namenode.NNStorage.getFinalizedEditsFileName;
@@ -168,14 +168,14 @@ public class TestNNStorageRetentionManag
StoragePurger mockPurger =
Mockito.mock(NNStorageRetentionManager.StoragePurger.class);
- ArgumentCaptor<FoundFSImage> imagesPurgedCaptor =
- ArgumentCaptor.forClass(FoundFSImage.class);
- ArgumentCaptor<FoundEditLog> logsPurgedCaptor =
- ArgumentCaptor.forClass(FoundEditLog.class);
+ ArgumentCaptor<FSImageFile> imagesPurgedCaptor =
+ ArgumentCaptor.forClass(FSImageFile.class);
+ ArgumentCaptor<EditLogFile> logsPurgedCaptor =
+ ArgumentCaptor.forClass(EditLogFile.class);
// Ask the manager to purge files we don't need any more
new NNStorageRetentionManager(conf,
- tc.mockStorage(), tc.mockEditLog(), mockPurger)
+ tc.mockStorage(), tc.mockEditLog(mockPurger), mockPurger)
.purgeOldStorage();
// Verify that it asked the purger to remove the correct files
@@ -186,7 +186,7 @@ public class TestNNStorageRetentionManag
// Check images
Set<String> purgedPaths = Sets.newHashSet();
- for (FoundFSImage purged : imagesPurgedCaptor.getAllValues()) {
+ for (FSImageFile purged : imagesPurgedCaptor.getAllValues()) {
purgedPaths.add(purged.getFile().toString());
}
Assert.assertEquals(Joiner.on(",").join(tc.expectedPurgedImages),
@@ -194,7 +194,7 @@ public class TestNNStorageRetentionManag
// Check images
purgedPaths.clear();
- for (FoundEditLog purged : logsPurgedCaptor.getAllValues()) {
+ for (EditLogFile purged : logsPurgedCaptor.getAllValues()) {
purgedPaths.add(purged.getFile().toString());
}
Assert.assertEquals(Joiner.on(",").join(tc.expectedPurgedLogs),
@@ -256,13 +256,14 @@ public class TestNNStorageRetentionManag
return mockStorageForDirs(sds.toArray(new StorageDirectory[0]));
}
- public FSEditLog mockEditLog() {
+ public FSEditLog mockEditLog(StoragePurger purger) {
final List<JournalManager> jms = Lists.newArrayList();
for (FakeRoot root : dirRoots.values()) {
if (!root.type.isOfType(NameNodeDirType.EDITS)) continue;
FileJournalManager fjm = new FileJournalManager(
root.mockStorageDir());
+ fjm.purger = purger;
jms.add(fjm);
}
@@ -272,17 +273,15 @@ public class TestNNStorageRetentionManag
@Override
public Void answer(InvocationOnMock invocation) throws Throwable {
Object[] args = invocation.getArguments();
- assert args.length == 2;
+ assert args.length == 1;
long txId = (Long) args[0];
- StoragePurger purger = (StoragePurger) args[1];
for (JournalManager jm : jms) {
- jm.purgeLogsOlderThan(txId, purger);
+ jm.purgeLogsOlderThan(txId);
}
return null;
}
- }).when(mockLog).purgeLogsOlderThan(
- Mockito.anyLong(), (StoragePurger) Mockito.anyObject());
+ }).when(mockLog).purgeLogsOlderThan(Mockito.anyLong());
return mockLog;
}
}