Skip to content

Commit 8333a39

Browse files
author
andrew mcdonald
committed
Back porting Avoids listing the sorted logs dir multiple times during log recovery. (#4874)
The log recovery code would list the sorted walog files multiple times during recovery. These changes modify the code to only list the files once. Also the listing is cached for a short period of time to improve the case of multiple tablet referencing the same walogs. This along with #4873 should result in much less traffic to the namenode when an entire accumulo cluster shutsdown and needs to recover. {"fundingSource": "41201", "team": "FED.ICGSA.OPS.MOE", "fshGit": "dummy-lo", "fshDocker": "sha256:20cf0045"}
1 parent 9bfc2c8 commit 8333a39

8 files changed

Lines changed: 237 additions & 51 deletions

File tree

‎server/tserver/src/main/java/org/apache/accumulo/tserver/TabletServer.java‎

Lines changed: 5 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1281,24 +1281,21 @@ public void minorCompactionStarted(CommitSession tablet, long lastUpdateSequence
12811281

12821282
public void recover(VolumeManager fs, KeyExtent extent, List<LogEntry> logEntries,
12831283
Set<String> tabletFiles, MutationReceiver mutationReceiver) throws IOException {
1284-
List<Path> recoveryDirs = new ArrayList<>();
12851284
List<LogEntry> sorted = new ArrayList<>(logEntries);
12861285
sorted.sort((e1, e2) -> (int) (e1.timestamp - e2.timestamp));
1286+
1287+
// Validate that recovery logs exist before attempting recovery
12871288
for (LogEntry entry : sorted) {
1288-
Path recovery = null;
12891289
Path finished = RecoveryPath.getRecoveryPath(new Path(entry.filename));
12901290
finished = SortedLogState.getFinishedMarkerPath(finished);
12911291
TabletServer.log.debug("Looking for " + finished);
1292-
if (fs.exists(finished)) {
1293-
recovery = finished.getParent();
1294-
}
1295-
if (recovery == null) {
1292+
if (!fs.exists(finished)) {
12961293
throw new IOException(
12971294
"Unable to find recovery files for extent " + extent + " logEntry: " + entry);
12981295
}
1299-
recoveryDirs.add(recovery);
13001296
}
1301-
logger.recover(getContext(), extent, recoveryDirs, tabletFiles, mutationReceiver);
1297+
1298+
logger.recover(getContext(), extent, sorted, tabletFiles, mutationReceiver);
13021299
}
13031300

13041301
public int createLogId() {

‎server/tserver/src/main/java/org/apache/accumulo/tserver/log/RecoveryLogsIterator.java‎

Lines changed: 10 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -66,16 +66,16 @@ public class RecoveryLogsIterator
6666
private final Iterator<Entry<Key,Value>> iter;
6767
private final CryptoEnvironment env = new CryptoEnvironmentImpl(CryptoEnvironment.Scope.RECOVERY);
6868

69-
public RecoveryLogsIterator(ServerContext context, List<Path> recoveryLogDirs, LogFileKey start,
70-
LogFileKey end, boolean checkFirstKey) throws IOException {
69+
public RecoveryLogsIterator(ServerContext context, List<ResolvedSortedLog> recoveryLogDirs,
70+
LogFileKey start, LogFileKey end, boolean checkFirstKey) throws IOException {
7171
this(context, recoveryLogDirs, start, end, checkFirstKey, null, null);
7272
}
7373

7474
/**
7575
* Scans the files in each recoveryLogDir over the range [start,end].
7676
*/
77-
public RecoveryLogsIterator(ServerContext context, List<Path> recoveryLogDirs, LogFileKey start,
78-
LogFileKey end, boolean checkFirstKey, Cache<String,Long> fileLenCache,
77+
public RecoveryLogsIterator(ServerContext context, List<ResolvedSortedLog> recoveryLogDirs,
78+
LogFileKey start, LogFileKey end, boolean checkFirstKey, Cache<String,Long> fileLenCache,
7979
CacheProvider cacheProvider) throws IOException {
8080

8181
List<Iterator<Entry<Key,Value>>> iterators = new ArrayList<>(recoveryLogDirs.size());
@@ -86,14 +86,15 @@ public RecoveryLogsIterator(ServerContext context, List<Path> recoveryLogDirs, L
8686
final CryptoService cryptoService = context.getCryptoFactory().getService(env,
8787
context.getConfiguration().getAllCryptoProperties());
8888

89-
for (Path logDir : recoveryLogDirs) {
90-
LOG.debug("Opening recovery log dir {}", logDir.getName());
91-
SortedSet<Path> logFiles = getFiles(vm, logDir);
92-
var fs = vm.getFileSystemByPath(logDir);
89+
for (ResolvedSortedLog logDir : recoveryLogDirs) {
90+
LOG.debug("Opening recovery log dir {}", logDir.getDir().getName());
91+
SortedSet<Path> logFiles = logDir.getChildren();
92+
var fs = vm.getFileSystemByPath(logDir.getDir());
9393

9494
// only check the first key once to prevent extra iterator creation and seeking
9595
if (checkFirstKey && !logFiles.isEmpty()) {
96-
validateFirstKey(context, cryptoService, fs, logFiles, logDir, fileLenCache, cacheProvider);
96+
validateFirstKey(context, cryptoService, fs, logFiles, logDir.getDir(), fileLenCache,
97+
cacheProvider);
9798
}
9899

99100
for (Path log : logFiles) {
Lines changed: 151 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,151 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* https://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing,
13+
* software distributed under the License is distributed on an
14+
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15+
* KIND, either express or implied. See the License for the
16+
* specific language governing permissions and limitations
17+
* under the License.
18+
*/
19+
package org.apache.accumulo.tserver.log;
20+
21+
import java.io.IOException;
22+
import java.util.Collections;
23+
import java.util.Comparator;
24+
import java.util.SortedSet;
25+
import java.util.TreeSet;
26+
27+
import org.apache.accumulo.core.tabletserver.log.LogEntry;
28+
import org.apache.accumulo.server.fs.VolumeManager;
29+
import org.apache.accumulo.server.log.SortedLogState;
30+
import org.apache.accumulo.server.manager.recovery.RecoveryPath;
31+
import org.apache.hadoop.fs.FileStatus;
32+
import org.apache.hadoop.fs.FileSystem;
33+
import org.apache.hadoop.fs.Path;
34+
35+
/**
36+
* Write ahead logs have two paths in DFS. There is the path of the original unsorted walog and the
37+
* path of the sorted walog. The purpose of this class is to convert the unsorted wal path to a
38+
* sorted wal path and validate the sorted dir exists and is finished.
39+
*/
40+
public class ResolvedSortedLog {
41+
42+
private final SortedSet<Path> children;
43+
private final LogEntry origin;
44+
private final Path sortedLogDir;
45+
46+
private ResolvedSortedLog(LogEntry origin, Path sortedLogDir, SortedSet<Path> children) {
47+
this.origin = origin;
48+
this.sortedLogDir = sortedLogDir;
49+
this.children = Collections.unmodifiableSortedSet(children);
50+
}
51+
52+
/**
53+
* @return the unsorted walog path from which this was created.
54+
*/
55+
public LogEntry getOrigin() {
56+
return origin;
57+
}
58+
59+
/**
60+
* @return the path of the directory in which sorted logs are stored
61+
*/
62+
public Path getDir() {
63+
return sortedLogDir;
64+
}
65+
66+
/**
67+
* @return When an unsorted walog is sorted the sorted data is stored in one or more rfiles, this
68+
* returns the paths of those rfiles.
69+
*/
70+
public SortedSet<Path> getChildren() {
71+
return children;
72+
}
73+
74+
@Override
75+
public String toString() {
76+
return sortedLogDir.toString();
77+
}
78+
79+
/**
80+
* For a given path of an unsorted walog check to see if the corresponding sorted log dir exists
81+
* and is finished. If it is return an immutable object containing information about the sorted
82+
* walogs.
83+
*/
84+
public static ResolvedSortedLog resolve(LogEntry logEntry, VolumeManager fs) throws IOException {
85+
86+
// convert the path of an unsorted log to the expected path for the corresponding sorted log
87+
// dir
88+
Path sortedLogPath = RecoveryPath.getRecoveryPath(new Path(logEntry.filename));
89+
90+
boolean foundFinish = false;
91+
// Path::getName compares the last component of each Path value. In this case, the last
92+
// component should
93+
// always have the format 'part-r-XXXXX.rf', where XXXXX are one-up values.
94+
SortedSet<Path> logFiles = new TreeSet<>(Comparator.comparing(Path::getName));
95+
for (FileStatus child : fs.listStatus(sortedLogPath)) {
96+
if (child.getPath().getName().startsWith("_")) {
97+
continue;
98+
}
99+
if (SortedLogState.isFinished(child.getPath().getName())) {
100+
foundFinish = true;
101+
continue;
102+
}
103+
if (SortedLogState.FAILED.getMarker().equals(child.getPath().getName())) {
104+
continue;
105+
}
106+
FileSystem ns = fs.getFileSystemByPath(child.getPath());
107+
Path fullLogPath = ns.makeQualified(child.getPath());
108+
logFiles.add(fullLogPath);
109+
}
110+
if (!foundFinish) {
111+
throw new IOException("Sort '" + SortedLogState.FINISHED.getMarker() + "' flag not found in "
112+
+ sortedLogPath + " for walog " + logEntry.filename);
113+
}
114+
115+
return new ResolvedSortedLog(logEntry, sortedLogPath, logFiles);
116+
}
117+
118+
/**
119+
* Create a ResolvedSortedLog directly from a sorted log directory path. This is useful for
120+
* diagnostic tools that operate directly on sorted recovery logs without going through the normal
121+
* recovery flow with LogEntry objects.
122+
*/
123+
public static ResolvedSortedLog fromSortedLogDir(Path sortedLogDir, VolumeManager fs)
124+
throws IOException {
125+
boolean foundFinish = false;
126+
SortedSet<Path> logFiles = new TreeSet<>(Comparator.comparing(Path::getName));
127+
for (FileStatus child : fs.listStatus(sortedLogDir)) {
128+
if (child.getPath().getName().startsWith("_")) {
129+
continue;
130+
}
131+
if (SortedLogState.isFinished(child.getPath().getName())) {
132+
foundFinish = true;
133+
continue;
134+
}
135+
if (SortedLogState.FAILED.getMarker().equals(child.getPath().getName())) {
136+
continue;
137+
}
138+
FileSystem ns = fs.getFileSystemByPath(child.getPath());
139+
Path fullLogPath = ns.makeQualified(child.getPath());
140+
logFiles.add(fullLogPath);
141+
}
142+
if (!foundFinish) {
143+
throw new IOException(
144+
"Sort '" + SortedLogState.FINISHED.getMarker() + "' flag not found in " + sortedLogDir);
145+
}
146+
147+
// Create a dummy LogEntry for the origin (used only for diagnostics)
148+
LogEntry dummyOrigin = new LogEntry(null, 0, sortedLogDir.toString());
149+
return new ResolvedSortedLog(dummyOrigin, sortedLogDir, logFiles);
150+
}
151+
}

‎server/tserver/src/main/java/org/apache/accumulo/tserver/log/SortedLogRecovery.java‎

Lines changed: 23 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -78,8 +78,10 @@ public SortedLogRecovery(ServerContext context, Cache<String,Long> fileLenCache,
7878
this.fileLenCache = fileLenCache;
7979
}
8080

81-
public boolean needsRecovery(KeyExtent extent, List<Path> recoveryDirs) throws IOException {
82-
Entry<Integer,List<Path>> maxEntry = findLogsThatDefineTablet(extent, recoveryDirs);
81+
public boolean needsRecovery(KeyExtent extent, List<ResolvedSortedLog> recoveryDirs)
82+
throws IOException {
83+
Entry<Integer,List<ResolvedSortedLog>> maxEntry =
84+
findLogsThatDefineTablet(extent, recoveryDirs);
8385
int tabletId = maxEntry.getKey();
8486
return tabletId != -1;
8587
}
@@ -115,7 +117,8 @@ static LogFileKey minKey(LogEvents event, int tabletId) {
115117
return key;
116118
}
117119

118-
private int findMaxTabletId(KeyExtent extent, List<Path> recoveryLogDirs) throws IOException {
120+
private int findMaxTabletId(KeyExtent extent, List<ResolvedSortedLog> recoveryLogDirs)
121+
throws IOException {
119122
int tabletId = -1;
120123

121124
try (var rli = new RecoveryLogsIterator(context, recoveryLogDirs, minKey(DEFINE_TABLET),
@@ -155,18 +158,18 @@ private int findMaxTabletId(KeyExtent extent, List<Path> recoveryLogDirs) throws
155158
* @return The maximum tablet ID observed AND the list of logs that contained the maximum tablet
156159
* ID.
157160
*/
158-
private Entry<Integer,List<Path>> findLogsThatDefineTablet(KeyExtent extent,
159-
List<Path> recoveryDirs) throws IOException {
160-
Map<Integer,List<Path>> logsThatDefineTablet = new HashMap<>();
161+
private Entry<Integer,List<ResolvedSortedLog>> findLogsThatDefineTablet(KeyExtent extent,
162+
List<ResolvedSortedLog> recoveryDirs) throws IOException {
163+
Map<Integer,List<ResolvedSortedLog>> logsThatDefineTablet = new HashMap<>();
161164

162-
for (Path walDir : recoveryDirs) {
165+
for (ResolvedSortedLog walDir : recoveryDirs) {
163166
int tabletId = findMaxTabletId(extent, Collections.singletonList(walDir));
164167
if (tabletId == -1) {
165-
log.debug("Did not find tablet {} in recovery log {}", extent, walDir.getName());
168+
log.debug("Did not find tablet {} in recovery log {}", extent, walDir.getDir().getName());
166169
} else {
167170
logsThatDefineTablet.computeIfAbsent(tabletId, k -> new ArrayList<>()).add(walDir);
168171
log.debug("Found tablet {} with id {} in recovery log {}", extent, tabletId,
169-
walDir.getName());
172+
walDir.getDir().getName());
170173
}
171174
}
172175

@@ -211,8 +214,8 @@ public Entry<LogFileKey,LogFileValue> next() {
211214

212215
}
213216

214-
private long findRecoverySeq(List<Path> recoveryLogs, Set<String> tabletFiles, int tabletId)
215-
throws IOException {
217+
private long findRecoverySeq(List<ResolvedSortedLog> recoveryLogs, Set<String> tabletFiles,
218+
int tabletId) throws IOException {
216219
HashSet<String> suffixes = new HashSet<>();
217220
for (String path : tabletFiles) {
218221
suffixes.add(getPathSuffix(path));
@@ -274,8 +277,8 @@ private long findRecoverySeq(List<Path> recoveryLogs, Set<String> tabletFiles, i
274277
return recoverySeq;
275278
}
276279

277-
private void playbackMutations(List<Path> recoveryLogs, MutationReceiver mr, int tabletId,
278-
long recoverySeq) throws IOException {
280+
private void playbackMutations(List<ResolvedSortedLog> recoveryLogs, MutationReceiver mr,
281+
int tabletId, long recoverySeq) throws IOException {
279282
LogFileKey start = minKey(MUTATION, tabletId);
280283
start.setSeq(recoverySeq);
281284

@@ -303,20 +306,21 @@ private void playbackMutations(List<Path> recoveryLogs, MutationReceiver mr, int
303306
}
304307
}
305308

306-
Collection<String> asNames(List<Path> recoveryLogs) {
307-
return Collections2.transform(recoveryLogs, Path::getName);
309+
Collection<String> asNames(List<ResolvedSortedLog> recoveryLogs) {
310+
return Collections2.transform(recoveryLogs, rl -> rl.getDir().getName());
308311
}
309312

310-
public void recover(KeyExtent extent, List<Path> recoveryDirs, Set<String> tabletFiles,
311-
MutationReceiver mr) throws IOException {
313+
public void recover(KeyExtent extent, List<ResolvedSortedLog> recoveryDirs,
314+
Set<String> tabletFiles, MutationReceiver mr) throws IOException {
312315

313-
Entry<Integer,List<Path>> maxEntry = findLogsThatDefineTablet(extent, recoveryDirs);
316+
Entry<Integer,List<ResolvedSortedLog>> maxEntry =
317+
findLogsThatDefineTablet(extent, recoveryDirs);
314318

315319
// A tablet may leave a tserver and then come back, in which case it would have a different and
316320
// higher tablet id. Only want to consider events in the log related to the last time the tablet
317321
// was loaded.
318322
int tabletId = maxEntry.getKey();
319-
List<Path> logsThatDefineTablet = maxEntry.getValue();
323+
List<ResolvedSortedLog> logsThatDefineTablet = maxEntry.getValue();
320324

321325
if (tabletId == -1) {
322326
log.info("Tablet {} is not defined in recovery logs {} ", extent, asNames(recoveryDirs));

‎server/tserver/src/main/java/org/apache/accumulo/tserver/log/TabletServerLogger.java‎

Lines changed: 20 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -65,6 +65,9 @@
6565
import org.slf4j.Logger;
6666
import org.slf4j.LoggerFactory;
6767

68+
import com.google.common.cache.Cache;
69+
import com.google.common.cache.CacheBuilder;
70+
6871
/**
6972
* Central logging facility for the TServerInfo.
7073
*
@@ -83,6 +86,9 @@ public class TabletServerLogger {
8386

8487
private final TabletServer tserver;
8588

89+
// Cache for resolved sorted logs to avoid repeated expensive I/O operations
90+
private final Cache<LogEntry,ResolvedSortedLog> sortedLogCache;
91+
8692
// The current logger
8793
private DfsLogger currentLog = null;
8894
private final SynchronousQueue<Object> nextLog = new SynchronousQueue<>();
@@ -164,6 +170,8 @@ public TabletServerLogger(TabletServer tserver, long maxSize, AtomicLong syncCou
164170
this.createRetry = null;
165171
this.writeRetryFactory = writeRetryFactory;
166172
this.maxAge = maxAge;
173+
this.sortedLogCache =
174+
CacheBuilder.newBuilder().expireAfterWrite(3, TimeUnit.SECONDS).maximumSize(1000).build();
167175
}
168176

169177
private DfsLogger initializeLoggers(final AtomicInteger logIdOut) throws IOException {
@@ -560,10 +568,17 @@ public long minorCompactionStarted(final CommitSession commitSession, final long
560568
return seq;
561569
}
562570

563-
private List<Path> resolve(Collection<LogEntry> walogs) {
564-
List<Path> sortedLogs = new ArrayList<>(walogs.size());
571+
private List<ResolvedSortedLog> resolve(Collection<LogEntry> walogs) throws IOException {
572+
List<ResolvedSortedLog> sortedLogs = new ArrayList<>(walogs.size());
573+
VolumeManager fs = tserver.getContext().getVolumeManager();
565574
for (var logEntry : walogs) {
566-
sortedLogs.add(new Path(logEntry.filename));
575+
try {
576+
ResolvedSortedLog resolvedLog =
577+
sortedLogCache.get(logEntry, () -> ResolvedSortedLog.resolve(logEntry, fs));
578+
sortedLogs.add(resolvedLog);
579+
} catch (Exception e) {
580+
throw new IOException("Failed to resolve sorted log for " + logEntry.filename, e);
581+
}
567582
}
568583
return sortedLogs;
569584
}
@@ -587,14 +602,14 @@ public boolean needsRecovery(ServerContext context, KeyExtent extent, Collection
587602
}
588603
}
589604

590-
public void recover(ServerContext context, KeyExtent extent, List<Path> recoveryDirs,
605+
public void recover(ServerContext context, KeyExtent extent, Collection<LogEntry> walogs,
591606
Set<String> tabletFiles, MutationReceiver mr) throws IOException {
592607
try {
593608
var resourceMgr = tserver.getResourceManager();
594609
var cacheProvider = createCacheProvider(resourceMgr);
595610
SortedLogRecovery recovery =
596611
new SortedLogRecovery(context, resourceMgr.getFileLenCache(), cacheProvider);
597-
recovery.recover(extent, recoveryDirs, tabletFiles, mr);
612+
recovery.recover(extent, resolve(walogs), tabletFiles, mr);
598613
} catch (Exception e) {
599614
throw new IOException(e);
600615
}

0 commit comments

Comments
 (0)