YARN-10189. Code cleanup in LeveldbRMStateStore. Contributed by Benjamin Teke
This commit is contained in:
parent
57c1e7ad9b
commit
f445487d50
|
@ -0,0 +1,131 @@
|
|||
/**
|
||||
* Licensed to the Apache Software Foundation (ASF) under one
|
||||
* or more contributor license agreements. See the NOTICE file
|
||||
* distributed with this work for additional information
|
||||
* regarding copyright ownership. The ASF licenses this file
|
||||
* to you under the Apache License, Version 2.0 (the
|
||||
* "License"); you may not use this file except in compliance
|
||||
* with the License. You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.apache.hadoop.yarn.server.resourcemanager;
|
||||
|
||||
import com.google.common.annotations.VisibleForTesting;
|
||||
import org.apache.hadoop.util.Time;
|
||||
import org.apache.hadoop.yarn.proto.YarnServerCommonProtos;
|
||||
import org.apache.hadoop.yarn.server.records.Version;
|
||||
import org.apache.hadoop.yarn.server.records.impl.pb.VersionPBImpl;
|
||||
import org.fusesource.leveldbjni.JniDBFactory;
|
||||
import org.fusesource.leveldbjni.internal.NativeDB;
|
||||
import org.iq80.leveldb.DB;
|
||||
import org.iq80.leveldb.DBException;
|
||||
import org.iq80.leveldb.Options;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.io.Closeable;
|
||||
import java.io.File;
|
||||
import java.io.IOException;
|
||||
import java.util.Timer;
|
||||
import java.util.TimerTask;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import static org.fusesource.leveldbjni.JniDBFactory.bytes;
|
||||
|
||||
public class DBManager implements Closeable {
|
||||
public static final Logger LOG =
|
||||
LoggerFactory.getLogger(DBManager.class);
|
||||
private DB db;
|
||||
private Timer compactionTimer;
|
||||
|
||||
public DB initDatabase(File configurationFile, Options options,
|
||||
Consumer<DB> initMethod) throws Exception {
|
||||
try {
|
||||
db = JniDBFactory.factory.open(configurationFile, options);
|
||||
} catch (NativeDB.DBException e) {
|
||||
if (e.isNotFound() || e.getMessage().contains(" does not exist ")) {
|
||||
LOG.info("Creating configuration version/database at {}",
|
||||
configurationFile);
|
||||
options.createIfMissing(true);
|
||||
try {
|
||||
db = JniDBFactory.factory.open(configurationFile, options);
|
||||
initMethod.accept(db);
|
||||
} catch (DBException dbErr) {
|
||||
throw new IOException(dbErr.getMessage(), dbErr);
|
||||
}
|
||||
} else {
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
|
||||
return db;
|
||||
}
|
||||
|
||||
public void close() throws IOException {
|
||||
if (compactionTimer != null) {
|
||||
compactionTimer.cancel();
|
||||
compactionTimer = null;
|
||||
}
|
||||
if (db != null) {
|
||||
db.close();
|
||||
db = null;
|
||||
}
|
||||
}
|
||||
|
||||
public void storeVersion(String versionKey, Version versionValue) {
|
||||
byte[] data = ((VersionPBImpl) versionValue).getProto().toByteArray();
|
||||
db.put(bytes(versionKey), data);
|
||||
}
|
||||
|
||||
public Version loadVersion(String versionKey) throws Exception {
|
||||
Version version = null;
|
||||
try {
|
||||
byte[] data = db.get(bytes(versionKey));
|
||||
if (data != null) {
|
||||
version = new VersionPBImpl(YarnServerCommonProtos.VersionProto
|
||||
.parseFrom(data));
|
||||
}
|
||||
} catch (DBException e) {
|
||||
throw new IOException(e);
|
||||
}
|
||||
return version;
|
||||
}
|
||||
|
||||
@VisibleForTesting
|
||||
public void setDb(DB db) {
|
||||
this.db = db;
|
||||
}
|
||||
|
||||
public void startCompactionTimer(long compactionIntervalMsec,
|
||||
String className) {
|
||||
if (compactionIntervalMsec > 0) {
|
||||
compactionTimer = new Timer(
|
||||
className + " compaction timer", true);
|
||||
compactionTimer.schedule(new CompactionTimerTask(),
|
||||
compactionIntervalMsec, compactionIntervalMsec);
|
||||
}
|
||||
}
|
||||
|
||||
private class CompactionTimerTask extends TimerTask {
|
||||
@Override
|
||||
public void run() {
|
||||
long start = Time.monotonicNow();
|
||||
LOG.info("Starting full compaction cycle");
|
||||
try {
|
||||
db.compactRange(null, null);
|
||||
} catch (DBException e) {
|
||||
LOG.error("Error compacting database", e);
|
||||
}
|
||||
long duration = Time.monotonicNow() - start;
|
||||
LOG.info("Full compaction cycle completed in " + duration + " msec");
|
||||
}
|
||||
}
|
||||
}
|
|
@ -29,9 +29,8 @@ import java.io.File;
|
|||
import java.io.IOException;
|
||||
import java.util.HashMap;
|
||||
import java.util.Map.Entry;
|
||||
import java.util.Timer;
|
||||
import java.util.TimerTask;
|
||||
|
||||
import org.apache.hadoop.yarn.server.resourcemanager.DBManager;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.apache.hadoop.conf.Configuration;
|
||||
|
@ -40,13 +39,11 @@ import org.apache.hadoop.fs.Path;
|
|||
import org.apache.hadoop.fs.permission.FsPermission;
|
||||
import org.apache.hadoop.io.IOUtils;
|
||||
import org.apache.hadoop.security.token.delegation.DelegationKey;
|
||||
import org.apache.hadoop.util.Time;
|
||||
import org.apache.hadoop.yarn.api.records.ApplicationAttemptId;
|
||||
import org.apache.hadoop.yarn.api.records.ApplicationId;
|
||||
import org.apache.hadoop.yarn.api.records.ReservationId;
|
||||
import org.apache.hadoop.yarn.conf.YarnConfiguration;
|
||||
import org.apache.hadoop.yarn.exceptions.YarnRuntimeException;
|
||||
import org.apache.hadoop.yarn.proto.YarnServerCommonProtos.VersionProto;
|
||||
import org.apache.hadoop.yarn.proto.YarnServerResourceManagerRecoveryProtos.AMRMTokenSecretManagerStateProto;
|
||||
import org.apache.hadoop.yarn.proto.YarnServerResourceManagerRecoveryProtos.EpochProto;
|
||||
import org.apache.hadoop.yarn.proto.YarnServerResourceManagerRecoveryProtos.ApplicationAttemptStateDataProto;
|
||||
|
@ -54,7 +51,6 @@ import org.apache.hadoop.yarn.proto.YarnServerResourceManagerRecoveryProtos.Appl
|
|||
import org.apache.hadoop.yarn.proto.YarnProtos.ReservationAllocationStateProto;
|
||||
import org.apache.hadoop.yarn.security.client.RMDelegationTokenIdentifier;
|
||||
import org.apache.hadoop.yarn.server.records.Version;
|
||||
import org.apache.hadoop.yarn.server.records.impl.pb.VersionPBImpl;
|
||||
import org.apache.hadoop.yarn.server.resourcemanager.recovery.records.AMRMTokenSecretManagerState;
|
||||
import org.apache.hadoop.yarn.server.resourcemanager.recovery.records.ApplicationAttemptStateData;
|
||||
import org.apache.hadoop.yarn.server.resourcemanager.recovery.records.ApplicationStateData;
|
||||
|
@ -64,8 +60,6 @@ import org.apache.hadoop.yarn.server.resourcemanager.recovery.records.impl.pb.AM
|
|||
import org.apache.hadoop.yarn.server.resourcemanager.recovery.records.impl.pb.ApplicationAttemptStateDataPBImpl;
|
||||
import org.apache.hadoop.yarn.server.resourcemanager.recovery.records.impl.pb.ApplicationStateDataPBImpl;
|
||||
import org.apache.hadoop.yarn.server.utils.LeveldbIterator;
|
||||
import org.fusesource.leveldbjni.JniDBFactory;
|
||||
import org.fusesource.leveldbjni.internal.NativeDB;
|
||||
import org.iq80.leveldb.DB;
|
||||
import org.iq80.leveldb.DBException;
|
||||
import org.iq80.leveldb.Options;
|
||||
|
@ -98,7 +92,7 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
.newInstance(1, 1);
|
||||
|
||||
private DB db;
|
||||
private Timer compactionTimer;
|
||||
private DBManager dbManager = new DBManager();
|
||||
private long compactionIntervalMsec;
|
||||
|
||||
private String getApplicationNodeKey(ApplicationId appId) {
|
||||
|
@ -130,7 +124,7 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
}
|
||||
|
||||
@Override
|
||||
protected void initInternal(Configuration conf) throws Exception {
|
||||
protected void initInternal(Configuration conf) {
|
||||
compactionIntervalMsec = conf.getLong(
|
||||
YarnConfiguration.RM_LEVELDB_COMPACTION_INTERVAL_SECS,
|
||||
YarnConfiguration.DEFAULT_RM_LEVELDB_COMPACTION_INTERVAL_SECS) * 1000;
|
||||
|
@ -155,55 +149,20 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
|
||||
@Override
|
||||
protected void startInternal() throws Exception {
|
||||
db = openDatabase();
|
||||
startCompactionTimer();
|
||||
}
|
||||
|
||||
protected DB openDatabase() throws Exception {
|
||||
Path storeRoot = createStorageDir();
|
||||
Options options = new Options();
|
||||
options.createIfMissing(false);
|
||||
LOG.info("Using state database at " + storeRoot + " for recovery");
|
||||
File dbfile = new File(storeRoot.toString());
|
||||
try {
|
||||
db = JniDBFactory.factory.open(dbfile, options);
|
||||
} catch (NativeDB.DBException e) {
|
||||
if (e.isNotFound() || e.getMessage().contains(" does not exist ")) {
|
||||
LOG.info("Creating state database at " + dbfile);
|
||||
options.createIfMissing(true);
|
||||
try {
|
||||
db = JniDBFactory.factory.open(dbfile, options);
|
||||
// store version
|
||||
storeVersion();
|
||||
} catch (DBException dbErr) {
|
||||
throw new IOException(dbErr.getMessage(), dbErr);
|
||||
}
|
||||
} else {
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
return db;
|
||||
}
|
||||
|
||||
private void startCompactionTimer() {
|
||||
if (compactionIntervalMsec > 0) {
|
||||
compactionTimer = new Timer(
|
||||
this.getClass().getSimpleName() + " compaction timer", true);
|
||||
compactionTimer.schedule(new CompactionTimerTask(),
|
||||
compactionIntervalMsec, compactionIntervalMsec);
|
||||
}
|
||||
db = dbManager.initDatabase(dbfile, options, (database) ->
|
||||
storeVersion(CURRENT_VERSION_INFO));
|
||||
dbManager.startCompactionTimer(compactionIntervalMsec,
|
||||
this.getClass().getSimpleName());
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void closeInternal() throws Exception {
|
||||
if (compactionTimer != null) {
|
||||
compactionTimer.cancel();
|
||||
compactionTimer = null;
|
||||
}
|
||||
if (db != null) {
|
||||
db.close();
|
||||
db = null;
|
||||
}
|
||||
dbManager.close();
|
||||
}
|
||||
|
||||
@VisibleForTesting
|
||||
|
@ -218,33 +177,22 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
|
||||
@Override
|
||||
protected Version loadVersion() throws Exception {
|
||||
Version version = null;
|
||||
try {
|
||||
byte[] data = db.get(bytes(VERSION_NODE));
|
||||
if (data != null) {
|
||||
version = new VersionPBImpl(VersionProto.parseFrom(data));
|
||||
}
|
||||
} catch (DBException e) {
|
||||
throw new IOException(e);
|
||||
}
|
||||
return version;
|
||||
return dbManager.loadVersion(VERSION_NODE);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void storeVersion() throws Exception {
|
||||
dbStoreVersion(CURRENT_VERSION_INFO);
|
||||
}
|
||||
|
||||
void dbStoreVersion(Version state) throws IOException {
|
||||
String key = VERSION_NODE;
|
||||
byte[] data = ((VersionPBImpl) state).getProto().toByteArray();
|
||||
try {
|
||||
db.put(bytes(key), data);
|
||||
storeVersion(CURRENT_VERSION_INFO);
|
||||
} catch (DBException e) {
|
||||
throw new IOException(e);
|
||||
}
|
||||
}
|
||||
|
||||
protected void storeVersion(Version version) {
|
||||
dbManager.storeVersion(VERSION_NODE, version);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected Version getCurrentVersion() {
|
||||
return CURRENT_VERSION_INFO;
|
||||
|
@ -279,9 +227,7 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
|
||||
private void loadReservationState(RMState rmState) throws IOException {
|
||||
int numReservations = 0;
|
||||
LeveldbIterator iter = null;
|
||||
try {
|
||||
iter = new LeveldbIterator(db);
|
||||
try (LeveldbIterator iter = new LeveldbIterator(db)) {
|
||||
iter.seek(bytes(RM_RESERVATION_KEY_PREFIX));
|
||||
while (iter.hasNext()) {
|
||||
Entry<byte[],byte[]> entry = iter.next();
|
||||
|
@ -313,10 +259,6 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
}
|
||||
} catch (DBException e) {
|
||||
throw new IOException(e);
|
||||
} finally {
|
||||
if (iter != null) {
|
||||
iter.close();
|
||||
}
|
||||
}
|
||||
LOG.info("Recovered " + numReservations + " reservations");
|
||||
}
|
||||
|
@ -331,9 +273,7 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
|
||||
private int loadRMDTSecretManagerKeys(RMState state) throws IOException {
|
||||
int numKeys = 0;
|
||||
LeveldbIterator iter = null;
|
||||
try {
|
||||
iter = new LeveldbIterator(db);
|
||||
try (LeveldbIterator iter = new LeveldbIterator(db)) {
|
||||
iter.seek(bytes(RM_DT_MASTER_KEY_KEY_PREFIX));
|
||||
while (iter.hasNext()) {
|
||||
Entry<byte[],byte[]> entry = iter.next();
|
||||
|
@ -352,10 +292,6 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
}
|
||||
} catch (DBException e) {
|
||||
throw new IOException(e);
|
||||
} finally {
|
||||
if (iter != null) {
|
||||
iter.close();
|
||||
}
|
||||
}
|
||||
return numKeys;
|
||||
}
|
||||
|
@ -373,9 +309,7 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
|
||||
private int loadRMDTSecretManagerTokens(RMState state) throws IOException {
|
||||
int numTokens = 0;
|
||||
LeveldbIterator iter = null;
|
||||
try {
|
||||
iter = new LeveldbIterator(db);
|
||||
try (LeveldbIterator iter = new LeveldbIterator(db)) {
|
||||
iter.seek(bytes(RM_DT_TOKEN_KEY_PREFIX));
|
||||
while (iter.hasNext()) {
|
||||
Entry<byte[],byte[]> entry = iter.next();
|
||||
|
@ -397,17 +331,13 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
}
|
||||
} catch (DBException e) {
|
||||
throw new IOException(e);
|
||||
} finally {
|
||||
if (iter != null) {
|
||||
iter.close();
|
||||
}
|
||||
}
|
||||
return numTokens;
|
||||
}
|
||||
|
||||
private RMDelegationTokenIdentifierData loadDelegationToken(byte[] data)
|
||||
throws IOException {
|
||||
RMDelegationTokenIdentifierData tokenData = null;
|
||||
RMDelegationTokenIdentifierData tokenData;
|
||||
DataInputStream in = new DataInputStream(new ByteArrayInputStream(data));
|
||||
try {
|
||||
tokenData = RMStateStoreUtils.readRMDelegationTokenIdentifierData(in);
|
||||
|
@ -419,7 +349,7 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
|
||||
private void loadRMDTSecretManagerTokenSequenceNumber(RMState state)
|
||||
throws IOException {
|
||||
byte[] data = null;
|
||||
byte[] data;
|
||||
try {
|
||||
data = db.get(bytes(RM_DT_SEQUENCE_NUMBER_KEY));
|
||||
} catch (DBException e) {
|
||||
|
@ -438,9 +368,7 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
private void loadRMApps(RMState state) throws IOException {
|
||||
int numApps = 0;
|
||||
int numAppAttempts = 0;
|
||||
LeveldbIterator iter = null;
|
||||
try {
|
||||
iter = new LeveldbIterator(db);
|
||||
try (LeveldbIterator iter = new LeveldbIterator(db)) {
|
||||
iter.seek(bytes(RM_APP_KEY_PREFIX));
|
||||
while (iter.hasNext()) {
|
||||
Entry<byte[],byte[]> entry = iter.next();
|
||||
|
@ -460,10 +388,6 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
}
|
||||
} catch (DBException e) {
|
||||
throw new IOException(e);
|
||||
} finally {
|
||||
if (iter != null) {
|
||||
iter.close();
|
||||
}
|
||||
}
|
||||
LOG.info("Recovered " + numApps + " applications and " + numAppAttempts
|
||||
+ " application attempts");
|
||||
|
@ -519,7 +443,7 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
@VisibleForTesting
|
||||
ApplicationStateData loadRMAppState(ApplicationId appId) throws IOException {
|
||||
String appKey = getApplicationNodeKey(appId);
|
||||
byte[] data = null;
|
||||
byte[] data;
|
||||
try {
|
||||
data = db.get(bytes(appKey));
|
||||
} catch (DBException e) {
|
||||
|
@ -535,7 +459,7 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
ApplicationAttemptStateData loadRMAppAttemptState(
|
||||
ApplicationAttemptId attemptId) throws IOException {
|
||||
String attemptKey = getApplicationAttemptNodeKey(attemptId);
|
||||
byte[] data = null;
|
||||
byte[] data;
|
||||
try {
|
||||
data = db.get(bytes(attemptKey));
|
||||
} catch (DBException e) {
|
||||
|
@ -643,8 +567,7 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
appState.getApplicationSubmissionContext().getApplicationId();
|
||||
String appKey = getApplicationNodeKey(appId);
|
||||
try {
|
||||
WriteBatch batch = db.createWriteBatch();
|
||||
try {
|
||||
try (WriteBatch batch = db.createWriteBatch()) {
|
||||
batch.delete(bytes(appKey));
|
||||
for (ApplicationAttemptId attemptId : appState.attempts.keySet()) {
|
||||
String attemptKey = getApplicationAttemptNodeKey(appKey, attemptId);
|
||||
|
@ -655,8 +578,6 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
+ appState.attempts.size() + " attempts" + " at " + appKey);
|
||||
}
|
||||
db.write(batch);
|
||||
} finally {
|
||||
batch.close();
|
||||
}
|
||||
} catch (DBException e) {
|
||||
throw new IOException(e);
|
||||
|
@ -668,8 +589,7 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
ReservationAllocationStateProto reservationAllocation, String planName,
|
||||
String reservationIdName) throws Exception {
|
||||
try {
|
||||
WriteBatch batch = db.createWriteBatch();
|
||||
try {
|
||||
try (WriteBatch batch = db.createWriteBatch()) {
|
||||
String key = getReservationNodeKey(planName, reservationIdName);
|
||||
if (LOG.isDebugEnabled()) {
|
||||
LOG.debug("Storing state for reservation " + reservationIdName
|
||||
|
@ -677,8 +597,6 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
}
|
||||
batch.put(bytes(key), reservationAllocation.toByteArray());
|
||||
db.write(batch);
|
||||
} finally {
|
||||
batch.close();
|
||||
}
|
||||
} catch (DBException e) {
|
||||
throw new IOException(e);
|
||||
|
@ -689,8 +607,7 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
protected void removeReservationState(String planName,
|
||||
String reservationIdName) throws Exception {
|
||||
try {
|
||||
WriteBatch batch = db.createWriteBatch();
|
||||
try {
|
||||
try (WriteBatch batch = db.createWriteBatch()) {
|
||||
String reservationKey =
|
||||
getReservationNodeKey(planName, reservationIdName);
|
||||
batch.delete(bytes(reservationKey));
|
||||
|
@ -699,8 +616,6 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
+ " plan " + planName + " at " + reservationKey);
|
||||
}
|
||||
db.write(batch);
|
||||
} finally {
|
||||
batch.close();
|
||||
}
|
||||
} catch (DBException e) {
|
||||
throw new IOException(e);
|
||||
|
@ -716,23 +631,20 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
LOG.debug("Storing token to " + tokenKey);
|
||||
}
|
||||
try {
|
||||
WriteBatch batch = db.createWriteBatch();
|
||||
try {
|
||||
try (WriteBatch batch = db.createWriteBatch()) {
|
||||
batch.put(bytes(tokenKey), tokenData.toByteArray());
|
||||
if(!isUpdate) {
|
||||
if (!isUpdate) {
|
||||
ByteArrayOutputStream bs = new ByteArrayOutputStream();
|
||||
try (DataOutputStream ds = new DataOutputStream(bs)) {
|
||||
ds.writeInt(tokenId.getSequenceNumber());
|
||||
}
|
||||
if (LOG.isDebugEnabled()) {
|
||||
LOG.debug("Storing " + tokenId.getSequenceNumber() + " to "
|
||||
+ RM_DT_SEQUENCE_NUMBER_KEY);
|
||||
+ RM_DT_SEQUENCE_NUMBER_KEY);
|
||||
}
|
||||
batch.put(bytes(RM_DT_SEQUENCE_NUMBER_KEY), bs.toByteArray());
|
||||
}
|
||||
db.write(batch);
|
||||
} finally {
|
||||
batch.close();
|
||||
}
|
||||
} catch (DBException e) {
|
||||
throw new IOException(e);
|
||||
|
@ -775,11 +687,8 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
LOG.debug("Storing token master key to " + dbKey);
|
||||
}
|
||||
ByteArrayOutputStream os = new ByteArrayOutputStream();
|
||||
DataOutputStream out = new DataOutputStream(os);
|
||||
try {
|
||||
try (DataOutputStream out = new DataOutputStream(os)) {
|
||||
masterKey.write(out);
|
||||
} finally {
|
||||
out.close();
|
||||
}
|
||||
try {
|
||||
db.put(bytes(dbKey), os.toByteArray());
|
||||
|
@ -836,9 +745,7 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
@VisibleForTesting
|
||||
int getNumEntriesInDatabase() throws IOException {
|
||||
int numEntries = 0;
|
||||
LeveldbIterator iter = null;
|
||||
try {
|
||||
iter = new LeveldbIterator(db);
|
||||
try (LeveldbIterator iter = new LeveldbIterator(db)) {
|
||||
iter.seekToFirst();
|
||||
while (iter.hasNext()) {
|
||||
Entry<byte[], byte[]> entry = iter.next();
|
||||
|
@ -847,26 +754,12 @@ public class LeveldbRMStateStore extends RMStateStore {
|
|||
}
|
||||
} catch (DBException e) {
|
||||
throw new IOException(e);
|
||||
} finally {
|
||||
if (iter != null) {
|
||||
iter.close();
|
||||
}
|
||||
}
|
||||
return numEntries;
|
||||
}
|
||||
|
||||
private class CompactionTimerTask extends TimerTask {
|
||||
@Override
|
||||
public void run() {
|
||||
long start = Time.monotonicNow();
|
||||
LOG.info("Starting full compaction cycle");
|
||||
try {
|
||||
db.compactRange(null, null);
|
||||
} catch (DBException e) {
|
||||
LOG.error("Error compacting database", e);
|
||||
}
|
||||
long duration = Time.monotonicNow() - start;
|
||||
LOG.info("Full compaction cycle completed in " + duration + " msec");
|
||||
}
|
||||
@VisibleForTesting
|
||||
protected void setDbManager(DBManager dbManager) {
|
||||
this.dbManager = dbManager;
|
||||
}
|
||||
}
|
||||
|
|
|
@ -19,20 +19,16 @@
|
|||
package org.apache.hadoop.yarn.server.resourcemanager.scheduler.capacity.conf;
|
||||
|
||||
import com.google.common.annotations.VisibleForTesting;
|
||||
import org.apache.hadoop.yarn.server.resourcemanager.DBManager;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.apache.hadoop.conf.Configuration;
|
||||
import org.apache.hadoop.fs.FileSystem;
|
||||
import org.apache.hadoop.fs.Path;
|
||||
import org.apache.hadoop.fs.permission.FsPermission;
|
||||
import org.apache.hadoop.util.Time;
|
||||
import org.apache.hadoop.yarn.conf.YarnConfiguration;
|
||||
import org.apache.hadoop.yarn.proto.YarnServerCommonProtos;
|
||||
import org.apache.hadoop.yarn.server.records.Version;
|
||||
import org.apache.hadoop.yarn.server.records.impl.pb.VersionPBImpl;
|
||||
import org.apache.hadoop.yarn.server.resourcemanager.RMContext;
|
||||
import org.fusesource.leveldbjni.JniDBFactory;
|
||||
import org.fusesource.leveldbjni.internal.NativeDB;
|
||||
import org.iq80.leveldb.DB;
|
||||
import org.iq80.leveldb.DBComparator;
|
||||
import org.iq80.leveldb.DBException;
|
||||
|
@ -52,9 +48,6 @@ import java.nio.charset.StandardCharsets;
|
|||
import java.util.LinkedList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Timer;
|
||||
import java.util.TimerTask;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import static org.fusesource.leveldbjni.JniDBFactory.bytes;
|
||||
|
||||
|
@ -73,6 +66,8 @@ public class LeveldbConfigurationStore extends YarnConfigurationStore {
|
|||
private static final String CONF_VERSION_KEY = "conf-version";
|
||||
|
||||
private DB db;
|
||||
private DBManager dbManager;
|
||||
private DBManager versionDbManager;
|
||||
private DB versionDb;
|
||||
private long maxLogs;
|
||||
private Configuration conf;
|
||||
|
@ -81,23 +76,25 @@ public class LeveldbConfigurationStore extends YarnConfigurationStore {
|
|||
@VisibleForTesting
|
||||
protected static final Version CURRENT_VERSION_INFO = Version
|
||||
.newInstance(0, 1);
|
||||
private long compactionIntervalMsec;
|
||||
|
||||
@Override
|
||||
public void initialize(Configuration config, Configuration schedConf,
|
||||
RMContext rmContext) throws IOException {
|
||||
this.conf = config;
|
||||
this.initSchedConf = schedConf;
|
||||
this.dbManager = new DBManager();
|
||||
this.versionDbManager = new DBManager();
|
||||
try {
|
||||
initDatabase();
|
||||
this.maxLogs = config.getLong(
|
||||
YarnConfiguration.RM_SCHEDCONF_MAX_LOGS,
|
||||
YarnConfiguration.DEFAULT_RM_SCHEDCONF_LEVELDB_MAX_LOGS);
|
||||
this.compactionIntervalMsec = config.getLong(
|
||||
long compactionIntervalMsec = config.getLong(
|
||||
YarnConfiguration.RM_SCHEDCONF_LEVELDB_COMPACTION_INTERVAL_SECS,
|
||||
YarnConfiguration
|
||||
.DEFAULT_RM_SCHEDCONF_LEVELDB_COMPACTION_INTERVAL_SECS) * 1000;
|
||||
startCompactionTimer();
|
||||
dbManager.startCompactionTimer(compactionIntervalMsec,
|
||||
this.getClass().getSimpleName());
|
||||
} catch (Exception e) {
|
||||
throw new IOException(e);
|
||||
}
|
||||
|
@ -116,7 +113,7 @@ public class LeveldbConfigurationStore extends YarnConfigurationStore {
|
|||
confOptions.createIfMissing(false);
|
||||
File confVersionFile = new File(confVersion.toString());
|
||||
|
||||
versionDb = initDatabaseHelper(confVersionFile, confOptions,
|
||||
versionDb = versionDbManager.initDatabase(confVersionFile, confOptions,
|
||||
this::initVersionDb);
|
||||
|
||||
Path storeRoot = createStorageDir(DB_NAME);
|
||||
|
@ -156,7 +153,7 @@ public class LeveldbConfigurationStore extends YarnConfigurationStore {
|
|||
});
|
||||
LOG.info("Using conf database at {}", storeRoot);
|
||||
File dbFile = new File(storeRoot.toString());
|
||||
db = initDatabaseHelper(dbFile, options, this::initDb);
|
||||
db = dbManager.initDatabase(dbFile, options, this::initDb);
|
||||
}
|
||||
|
||||
private void initVersionDb(DB database) {
|
||||
|
@ -172,30 +169,6 @@ public class LeveldbConfigurationStore extends YarnConfigurationStore {
|
|||
increaseConfigVersion();
|
||||
}
|
||||
|
||||
private DB initDatabaseHelper(File configurationFile, Options options,
|
||||
Consumer<DB> initMethod) throws Exception {
|
||||
DB database;
|
||||
try {
|
||||
database = JniDBFactory.factory.open(configurationFile, options);
|
||||
} catch (NativeDB.DBException e) {
|
||||
if (e.isNotFound() || e.getMessage().contains(" does not exist ")) {
|
||||
LOG.info("Creating configuration version/database at {}",
|
||||
configurationFile);
|
||||
options.createIfMissing(true);
|
||||
try {
|
||||
database = JniDBFactory.factory.open(configurationFile, options);
|
||||
initMethod.accept(database);
|
||||
} catch (DBException dbErr) {
|
||||
throw new IOException(dbErr.getMessage(), dbErr);
|
||||
}
|
||||
} else {
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
|
||||
return database;
|
||||
}
|
||||
|
||||
private Path createStorageDir(String storageName) throws IOException {
|
||||
Path root = getStorageDir(storageName);
|
||||
FileSystem fs = FileSystem.getLocal(conf);
|
||||
|
@ -214,12 +187,8 @@ public class LeveldbConfigurationStore extends YarnConfigurationStore {
|
|||
|
||||
@Override
|
||||
public void close() throws IOException {
|
||||
if (db != null) {
|
||||
db.close();
|
||||
}
|
||||
if (versionDb != null) {
|
||||
versionDb.close();
|
||||
}
|
||||
dbManager.close();
|
||||
versionDbManager.close();
|
||||
}
|
||||
|
||||
@Override
|
||||
|
@ -313,28 +282,9 @@ public class LeveldbConfigurationStore extends YarnConfigurationStore {
|
|||
return null; // unimplemented
|
||||
}
|
||||
|
||||
private void startCompactionTimer() {
|
||||
if (compactionIntervalMsec > 0) {
|
||||
Timer compactionTimer = new Timer(
|
||||
this.getClass().getSimpleName() + " compaction timer", true);
|
||||
compactionTimer.schedule(new CompactionTimerTask(),
|
||||
compactionIntervalMsec, compactionIntervalMsec);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public Version getConfStoreVersion() throws Exception {
|
||||
Version version = null;
|
||||
try {
|
||||
byte[] data = db.get(bytes(VERSION_KEY));
|
||||
if (data != null) {
|
||||
version = new VersionPBImpl(YarnServerCommonProtos.VersionProto
|
||||
.parseFrom(data));
|
||||
}
|
||||
} catch (DBException e) {
|
||||
throw new IOException(e);
|
||||
}
|
||||
return version;
|
||||
return dbManager.loadVersion(VERSION_KEY);
|
||||
}
|
||||
|
||||
@VisibleForTesting
|
||||
|
@ -345,37 +295,20 @@ public class LeveldbConfigurationStore extends YarnConfigurationStore {
|
|||
|
||||
@Override
|
||||
public void storeVersion() throws Exception {
|
||||
storeVersion(CURRENT_VERSION_INFO);
|
||||
}
|
||||
|
||||
@VisibleForTesting
|
||||
protected void storeVersion(Version version) throws Exception {
|
||||
byte[] data = ((VersionPBImpl) version).getProto()
|
||||
.toByteArray();
|
||||
try {
|
||||
db.put(bytes(VERSION_KEY), data);
|
||||
storeVersion(CURRENT_VERSION_INFO);
|
||||
} catch (DBException e) {
|
||||
throw new IOException(e);
|
||||
}
|
||||
}
|
||||
|
||||
@VisibleForTesting
|
||||
protected void storeVersion(Version version) {
|
||||
dbManager.storeVersion(VERSION_KEY, version);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Version getCurrentVersion() {
|
||||
return CURRENT_VERSION_INFO;
|
||||
}
|
||||
|
||||
private class CompactionTimerTask extends TimerTask {
|
||||
@Override
|
||||
public void run() {
|
||||
long start = Time.monotonicNow();
|
||||
LOG.info("Starting full compaction cycle");
|
||||
try {
|
||||
db.compactRange(null, null);
|
||||
} catch (DBException e) {
|
||||
LOG.error("Error compacting database", e);
|
||||
}
|
||||
long duration = Time.monotonicNow() - start;
|
||||
LOG.info("Full compaction cycle completed in {} msec", duration);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
@ -25,14 +25,17 @@ import static org.mockito.Mockito.verify;
|
|||
|
||||
import java.io.File;
|
||||
import java.io.IOException;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import org.apache.hadoop.fs.FileUtil;
|
||||
import org.apache.hadoop.yarn.conf.YarnConfiguration;
|
||||
import org.apache.hadoop.yarn.server.records.Version;
|
||||
import org.apache.hadoop.yarn.server.resourcemanager.DBManager;
|
||||
import org.apache.hadoop.yarn.server.resourcemanager.rmapp.RMApp;
|
||||
import org.apache.hadoop.yarn.server.resourcemanager.rmapp.attempt.RMAppAttempt;
|
||||
import org.fusesource.leveldbjni.JniDBFactory;
|
||||
import org.iq80.leveldb.DB;
|
||||
import org.iq80.leveldb.Options;
|
||||
import org.junit.After;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
@ -125,15 +128,19 @@ public class TestLeveldbRMStateStore extends RMStateStoreTestBase {
|
|||
}
|
||||
|
||||
@Test(timeout = 60000)
|
||||
public void testCompactionCycle() throws Exception {
|
||||
public void testCompactionCycle() {
|
||||
final DB mockdb = mock(DB.class);
|
||||
conf.setLong(YarnConfiguration.RM_LEVELDB_COMPACTION_INTERVAL_SECS, 1);
|
||||
stateStore = new LeveldbRMStateStore() {
|
||||
stateStore = new LeveldbRMStateStore();
|
||||
DBManager dbManager = new DBManager() {
|
||||
@Override
|
||||
protected DB openDatabase() throws Exception {
|
||||
public DB initDatabase(File configurationFile, Options options,
|
||||
Consumer<DB> initMethod) {
|
||||
return mockdb;
|
||||
}
|
||||
};
|
||||
dbManager.setDb(mockdb);
|
||||
stateStore.setDbManager(dbManager);
|
||||
stateStore.init(conf);
|
||||
stateStore.start();
|
||||
verify(mockdb, timeout(10000)).compactRange(
|
||||
|
@ -175,12 +182,12 @@ public class TestLeveldbRMStateStore extends RMStateStoreTestBase {
|
|||
}
|
||||
|
||||
@Override
|
||||
public void writeVersion(Version version) throws Exception {
|
||||
stateStore.dbStoreVersion(version);
|
||||
public void writeVersion(Version version) {
|
||||
stateStore.storeVersion(version);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Version getCurrentVersion() throws Exception {
|
||||
public Version getCurrentVersion() {
|
||||
return stateStore.getCurrentVersion();
|
||||
}
|
||||
|
||||
|
|
Loading…
Reference in New Issue