Author: todd
Date: Thu Sep 6 07:03:57 2012
New Revision: 1381482
URL: http://svn.apache.org/viewvc?rev=1381482&view=rev
Log:
HDFS-3726. If a logger misses an RPC, don't retry that logger until next
segment. Contributed by Todd Lipcon.
Modified:
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/CHANGES.HDFS-3077.txt
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/qjournal/client/IPCLoggerChannel.java
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/qjournal/server/Journal.java
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/client/TestIPCLoggerChannel.java
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/server/TestJournal.java
Modified:
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/CHANGES.HDFS-3077.txt
URL:
http://svn.apache.org/viewvc/hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/CHANGES.HDFS-3077.txt?rev=1381482&r1=1381481&r2=1381482&view=diff
==============================================================================
---
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/CHANGES.HDFS-3077.txt
(original)
+++
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/CHANGES.HDFS-3077.txt
Thu Sep 6 07:03:57 2012
@@ -46,3 +46,5 @@ HDFS-3884. Journal format() should reset
HDFS-3870. Add metrics to JournalNode (todd)
HDFS-3891. Make selectInputStreams throw IOE instead of RTE (todd)
+
+HDFS-3726. If a logger misses an RPC, don't retry that logger until next
segment (todd)
Modified:
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/qjournal/client/IPCLoggerChannel.java
URL:
http://svn.apache.org/viewvc/hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/qjournal/client/IPCLoggerChannel.java?rev=1381482&r1=1381481&r2=1381482&view=diff
==============================================================================
---
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/qjournal/client/IPCLoggerChannel.java
(original)
+++
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/qjournal/client/IPCLoggerChannel.java
Thu Sep 6 07:03:57 2012
@@ -25,13 +25,13 @@ import java.util.concurrent.Callable;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
-import java.util.concurrent.ScheduledExecutorService;
import org.apache.hadoop.classification.InterfaceAudience;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.hdfs.DFSConfigKeys;
import org.apache.hadoop.hdfs.protocol.HdfsConstants;
import org.apache.hadoop.hdfs.protocolPB.PBHelper;
+import org.apache.hadoop.hdfs.qjournal.protocol.JournalOutOfSyncException;
import org.apache.hadoop.hdfs.qjournal.protocol.QJournalProtocol;
import
org.apache.hadoop.hdfs.qjournal.protocol.QJournalProtocolProtos.GetEditLogManifestResponseProto;
import
org.apache.hadoop.hdfs.qjournal.protocol.QJournalProtocolProtos.GetJournalStateResponseProto;
@@ -98,6 +98,13 @@ public class IPCLoggerChannel implements
*/
private final int queueSizeLimitBytes;
+ /**
+ * If this logger misses some edits, or restarts in the middle of
+ * a segment, the writer won't be able to write any more edits until
+ * the beginning of the next segment. Upon detecting this situation,
+ * the writer sets this flag to true to avoid sending useless RPCs.
+ */
+ private boolean outOfSync = false;
static final Factory FACTORY = new AsyncLogger.Factory() {
@Override
@@ -212,6 +219,15 @@ public class IPCLoggerChannel implements
public synchronized int getQueuedEditsSize() {
return queuedEditsSizeBytes;
}
+
+ /**
+ * @return true if the server has gotten out of sync from the client,
+ * and thus a log roll is required for this logger to successfully start
+ * logging more edits.
+ */
+ public synchronized boolean isOutOfSync() {
+ return outOfSync;
+ }
@VisibleForTesting
void waitForAllPendingCalls() throws InterruptedException {
@@ -265,8 +281,22 @@ public class IPCLoggerChannel implements
ret = executor.submit(new Callable<Void>() {
@Override
public Void call() throws IOException {
- getProxy().journal(createReqInfo(),
- segmentTxId, firstTxnId, numTxns, data);
+ throwIfOutOfSync();
+
+ try {
+ getProxy().journal(createReqInfo(),
+ segmentTxId, firstTxnId, numTxns, data);
+ } catch (IOException e) {
+ QuorumJournalManager.LOG.warn(
+ "Remote journal " + IPCLoggerChannel.this + " failed to " +
+ "write txns " + firstTxnId + "-" + (firstTxnId + numTxns - 1) +
+ ". Will try to write to this JN again after the next " +
+ "log roll.", e);
+ synchronized (IPCLoggerChannel.this) {
+ outOfSync = true;
+ }
+ throw e;
+ }
synchronized (IPCLoggerChannel.this) {
highestAckedTxId = firstTxnId + numTxns - 1;
}
@@ -298,6 +328,15 @@ public class IPCLoggerChannel implements
return ret;
}
+ private synchronized void throwIfOutOfSync() throws
JournalOutOfSyncException {
+ if (outOfSync) {
+ // TODO: send a "heartbeat" here so that the remote node knows the newest
+ // committed txid, for metrics purposes
+ throw new JournalOutOfSyncException(
+ "Journal disabled until next roll");
+ }
+ }
+
private synchronized void reserveQueueSpace(int size)
throws LoggerTooFarBehindException {
Preconditions.checkArgument(size >= 0);
@@ -330,6 +369,15 @@ public class IPCLoggerChannel implements
@Override
public Void call() throws IOException {
getProxy().startLogSegment(createReqInfo(), txid);
+ synchronized (IPCLoggerChannel.this) {
+ if (outOfSync) {
+ outOfSync = false;
+ QuorumJournalManager.LOG.info(
+ "Restarting previously-stopped writes to " +
+ IPCLoggerChannel.this + " in segment starting at txid " +
+ txid);
+ }
+ }
return null;
}
});
@@ -341,6 +389,8 @@ public class IPCLoggerChannel implements
return executor.submit(new Callable<Void>() {
@Override
public Void call() throws IOException {
+ throwIfOutOfSync();
+
getProxy().finalizeLogSegment(createReqInfo(), startTxId, endTxId);
return null;
}
@@ -415,5 +465,8 @@ public class IPCLoggerChannel implements
if (behind > 0) {
sb.append(" (" + behind + " behind)");
}
+ if (outOfSync) {
+ sb.append(" (will re-join on next segment)");
+ }
}
-}
\ No newline at end of file
+}
Modified:
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/qjournal/server/Journal.java
URL:
http://svn.apache.org/viewvc/hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/qjournal/server/Journal.java?rev=1381482&r1=1381481&r2=1381482&view=diff
==============================================================================
---
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/qjournal/server/Journal.java
(original)
+++
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/qjournal/server/Journal.java
Thu Sep 6 07:03:57 2012
@@ -33,6 +33,7 @@ import org.apache.hadoop.fs.FileUtil;
import org.apache.hadoop.hdfs.protocol.HdfsConstants;
import org.apache.hadoop.hdfs.qjournal.protocol.JournalNotFormattedException;
import org.apache.hadoop.hdfs.qjournal.protocol.QJournalProtocol;
+import org.apache.hadoop.hdfs.qjournal.protocol.JournalOutOfSyncException;
import
org.apache.hadoop.hdfs.qjournal.protocol.QJournalProtocolProtos.NewEpochResponseProto;
import
org.apache.hadoop.hdfs.qjournal.protocol.QJournalProtocolProtos.PersistedRecoveryPaxosData;
import
org.apache.hadoop.hdfs.qjournal.protocol.QJournalProtocolProtos.PrepareRecoveryResponseProto;
@@ -273,14 +274,8 @@ class Journal implements Closeable {
int numTxns, byte[] records) throws IOException {
checkFormatted();
checkWriteRequest(reqInfo);
-
- // TODO: if a JN goes down and comes back up, then it will throw
- // this exception on every edit. We should instead send back
- // a response indicating the log needs to be rolled, which would
- // mark the logger on the client side as "pending" -- and have the
- // NN code look for this condition and trigger a roll when it happens.
- // That way the node can catch back up and rejoin
- Preconditions.checkState(curSegment != null,
+
+ checkSync(curSegment != null,
"Can't write, no segment open");
if (curSegmentTxId != segmentTxId) {
@@ -297,7 +292,7 @@ class Journal implements Closeable {
+ " but current segment is " + curSegmentTxId);
}
- Preconditions.checkState(nextTxId == firstTxnId,
+ checkSync(nextTxId == firstTxnId,
"Can't write txid " + firstTxnId + " expecting nextTxId=" + nextTxId);
long lastTxnId = firstTxnId + numTxns - 1;
@@ -379,6 +374,18 @@ class Journal implements Closeable {
}
/**
+ * @throws JournalOutOfSyncException if the given expression is not true.
+ * The message of the exception is formatted using the 'msg' and
+ * 'formatArgs' parameters.
+ */
+ private void checkSync(boolean expression, String msg,
+ Object... formatArgs) throws JournalOutOfSyncException {
+ if (!expression) {
+ throw new JournalOutOfSyncException(String.format(msg, formatArgs));
+ }
+ }
+
+ /**
* Start a new segment at the given txid. The previous segment
* must have already been finalized.
*/
@@ -453,7 +460,7 @@ class Journal implements Closeable {
FileJournalManager.EditLogFile elf = fjm.getLogFile(startTxId);
if (elf == null) {
- throw new IllegalStateException("No log file to finalize at " +
+ throw new JournalOutOfSyncException("No log file to finalize at " +
"transaction ID " + startTxId);
}
@@ -463,8 +470,8 @@ class Journal implements Closeable {
LOG.info("Validating log about to be finalized: " + elf);
elf.validateLog();
-
- Preconditions.checkState(elf.getLastTxId() == endTxId,
+
+ checkSync(elf.getLastTxId() == endTxId,
"Trying to finalize log %s-%s, but current state of log " +
"is %s", startTxId, endTxId, elf);
fjm.finalizeLogSegment(startTxId, endTxId);
Modified:
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/client/TestIPCLoggerChannel.java
URL:
http://svn.apache.org/viewvc/hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/client/TestIPCLoggerChannel.java?rev=1381482&r1=1381481&r2=1381482&view=diff
==============================================================================
---
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/client/TestIPCLoggerChannel.java
(original)
+++
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/client/TestIPCLoggerChannel.java
Thu Sep 6 07:03:57 2012
@@ -128,5 +128,51 @@ public class TestIPCLoggerChannel {
}
}, 10, 1000);
}
+
+ /**
+ * Test that, if the remote node gets unsynchronized (eg some edits were
+ * missed or the node rebooted), the client stops sending edits until
+ * the next roll. Test for HDFS-3726.
+ */
+ @Test
+ public void testStopSendingEditsWhenOutOfSync() throws Exception {
+ Mockito.doThrow(new IOException("injected error"))
+ .when(mockProxy).journal(
+ Mockito.<RequestInfo>any(),
+ Mockito.eq(1L), Mockito.eq(1L),
+ Mockito.eq(1), Mockito.same(FAKE_DATA));
+ try {
+ ch.sendEdits(1L, 1L, 1, FAKE_DATA).get();
+ fail("Injected JOOSE did not cause sendEdits() to throw");
+ } catch (ExecutionException ee) {
+ GenericTestUtils.assertExceptionContains("injected", ee);
+ }
+ Mockito.verify(mockProxy).journal(
+ Mockito.<RequestInfo>any(),
+ Mockito.eq(1L), Mockito.eq(1L),
+ Mockito.eq(1), Mockito.same(FAKE_DATA));
+
+ assertTrue(ch.isOutOfSync());
+
+ try {
+ ch.sendEdits(1L, 2L, 1, FAKE_DATA).get();
+ fail("sendEdits() should throw until next roll");
+ } catch (ExecutionException ee) {
+ GenericTestUtils.assertExceptionContains("disabled until next roll",
+ ee.getCause());
+ }
+
+ // It should have failed without even sending an RPC, since it was not
sync.
+ Mockito.verify(mockProxy, Mockito.never()).journal(
+ Mockito.<RequestInfo>any(),
+ Mockito.eq(1L), Mockito.eq(2L),
+ Mockito.eq(1), Mockito.same(FAKE_DATA));
+
+ // After a roll, sending new edits should not fail.
+ ch.startLogSegment(3L).get();
+ assertFalse(ch.isOutOfSync());
+
+ ch.sendEdits(3L, 3L, 1, FAKE_DATA).get();
+ }
}
Modified:
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/server/TestJournal.java
URL:
http://svn.apache.org/viewvc/hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/server/TestJournal.java?rev=1381482&r1=1381481&r2=1381482&view=diff
==============================================================================
---
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/server/TestJournal.java
(original)
+++
hadoop/common/branches/HDFS-3077/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/qjournal/server/TestJournal.java
Thu Sep 6 07:03:57 2012
@@ -26,6 +26,7 @@ import org.apache.hadoop.fs.FileUtil;
import org.apache.hadoop.hdfs.MiniDFSCluster;
import org.apache.hadoop.hdfs.qjournal.QJMTestUtil;
import org.apache.hadoop.hdfs.qjournal.protocol.RequestInfo;
+import org.apache.hadoop.hdfs.qjournal.protocol.JournalOutOfSyncException;
import
org.apache.hadoop.hdfs.qjournal.protocol.QJournalProtocolProtos.NewEpochResponseProto;
import
org.apache.hadoop.hdfs.qjournal.protocol.QJournalProtocolProtos.NewEpochResponseProtoOrBuilder;
import org.apache.hadoop.hdfs.qjournal.server.Journal;
@@ -219,9 +220,9 @@ public class TestJournal {
try {
journal.finalizeLogSegment(makeRI(3), 1, 6);
fail("did not fail to finalize");
- } catch (IllegalStateException ise) {
+ } catch (JournalOutOfSyncException e) {
GenericTestUtils.assertExceptionContains(
- "but current state of log is", ise);
+ "but current state of log is", e);
}
// Check that, even if we re-construct the journal by scanning the
@@ -232,9 +233,9 @@ public class TestJournal {
try {
journal.finalizeLogSegment(makeRI(4), 1, 6);
fail("did not fail to finalize");
- } catch (IllegalStateException ise) {
+ } catch (JournalOutOfSyncException e) {
GenericTestUtils.assertExceptionContains(
- "but current state of log is", ise);
+ "but current state of log is", e);
}
}
@@ -248,9 +249,9 @@ public class TestJournal {
try {
journal.finalizeLogSegment(makeRI(1), 1000, 1001);
fail("did not fail to finalize");
- } catch (IllegalStateException ise) {
+ } catch (JournalOutOfSyncException e) {
GenericTestUtils.assertExceptionContains(
- "No log file to finalize at transaction ID 1000", ise);
+ "No log file to finalize at transaction ID 1000", e);
}
}