HBASE-12533: staging directories are not deleted after secure bulk load
This commit is contained in:
parent
def45c2185
commit
643de2df04
|
@ -26,6 +26,7 @@ import org.apache.hadoop.hbase.util.ByteStringer;
|
||||||
|
|
||||||
import org.apache.hadoop.hbase.classification.InterfaceAudience;
|
import org.apache.hadoop.hbase.classification.InterfaceAudience;
|
||||||
import org.apache.hadoop.fs.Path;
|
import org.apache.hadoop.fs.Path;
|
||||||
|
import org.apache.hadoop.hbase.HConstants;
|
||||||
import org.apache.hadoop.hbase.TableName;
|
import org.apache.hadoop.hbase.TableName;
|
||||||
import org.apache.hadoop.hbase.ipc.BlockingRpcCallback;
|
import org.apache.hadoop.hbase.ipc.BlockingRpcCallback;
|
||||||
import org.apache.hadoop.hbase.ipc.CoprocessorRpcChannel;
|
import org.apache.hadoop.hbase.ipc.CoprocessorRpcChannel;
|
||||||
|
@ -55,13 +56,10 @@ public class SecureBulkLoadClient {
|
||||||
|
|
||||||
public String prepareBulkLoad(final TableName tableName) throws IOException {
|
public String prepareBulkLoad(final TableName tableName) throws IOException {
|
||||||
try {
|
try {
|
||||||
return
|
CoprocessorRpcChannel channel = table.coprocessorService(HConstants.EMPTY_START_ROW);
|
||||||
table.coprocessorService(SecureBulkLoadProtos.SecureBulkLoadService.class,
|
SecureBulkLoadProtos.SecureBulkLoadService instance =
|
||||||
EMPTY_START_ROW,
|
ProtobufUtil.newServiceStub(SecureBulkLoadProtos.SecureBulkLoadService.class, channel);
|
||||||
LAST_ROW,
|
|
||||||
new Batch.Call<SecureBulkLoadProtos.SecureBulkLoadService,String>() {
|
|
||||||
@Override
|
|
||||||
public String call(SecureBulkLoadProtos.SecureBulkLoadService instance) throws IOException {
|
|
||||||
ServerRpcController controller = new ServerRpcController();
|
ServerRpcController controller = new ServerRpcController();
|
||||||
|
|
||||||
BlockingRpcCallback<SecureBulkLoadProtos.PrepareBulkLoadResponse> rpcCallback =
|
BlockingRpcCallback<SecureBulkLoadProtos.PrepareBulkLoadResponse> rpcCallback =
|
||||||
|
@ -79,9 +77,8 @@ public class SecureBulkLoadClient {
|
||||||
if (controller.failedOnException()) {
|
if (controller.failedOnException()) {
|
||||||
throw controller.getFailedOn();
|
throw controller.getFailedOn();
|
||||||
}
|
}
|
||||||
|
|
||||||
return response.getBulkToken();
|
return response.getBulkToken();
|
||||||
}
|
|
||||||
}).entrySet().iterator().next().getValue();
|
|
||||||
} catch (Throwable throwable) {
|
} catch (Throwable throwable) {
|
||||||
throw new IOException(throwable);
|
throw new IOException(throwable);
|
||||||
}
|
}
|
||||||
|
@ -89,13 +86,10 @@ public class SecureBulkLoadClient {
|
||||||
|
|
||||||
public void cleanupBulkLoad(final String bulkToken) throws IOException {
|
public void cleanupBulkLoad(final String bulkToken) throws IOException {
|
||||||
try {
|
try {
|
||||||
table.coprocessorService(SecureBulkLoadProtos.SecureBulkLoadService.class,
|
CoprocessorRpcChannel channel = table.coprocessorService(HConstants.EMPTY_START_ROW);
|
||||||
EMPTY_START_ROW,
|
SecureBulkLoadProtos.SecureBulkLoadService instance =
|
||||||
LAST_ROW,
|
ProtobufUtil.newServiceStub(SecureBulkLoadProtos.SecureBulkLoadService.class, channel);
|
||||||
new Batch.Call<SecureBulkLoadProtos.SecureBulkLoadService, String>() {
|
|
||||||
|
|
||||||
@Override
|
|
||||||
public String call(SecureBulkLoadProtos.SecureBulkLoadService instance) throws IOException {
|
|
||||||
ServerRpcController controller = new ServerRpcController();
|
ServerRpcController controller = new ServerRpcController();
|
||||||
|
|
||||||
BlockingRpcCallback<SecureBulkLoadProtos.CleanupBulkLoadResponse> rpcCallback =
|
BlockingRpcCallback<SecureBulkLoadProtos.CleanupBulkLoadResponse> rpcCallback =
|
||||||
|
@ -112,9 +106,6 @@ public class SecureBulkLoadClient {
|
||||||
if (controller.failedOnException()) {
|
if (controller.failedOnException()) {
|
||||||
throw controller.getFailedOn();
|
throw controller.getFailedOn();
|
||||||
}
|
}
|
||||||
return null;
|
|
||||||
}
|
|
||||||
});
|
|
||||||
} catch (Throwable throwable) {
|
} catch (Throwable throwable) {
|
||||||
throw new IOException(throwable);
|
throw new IOException(throwable);
|
||||||
}
|
}
|
||||||
|
|
|
@ -198,10 +198,7 @@ public class SecureBulkLoadEndpoint extends SecureBulkLoadService
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fs.delete(createStagingDir(baseStagingDir,
|
fs.delete(new Path(request.getBulkToken()), true);
|
||||||
getActiveUser(),
|
|
||||||
new Path(request.getBulkToken()).getName()),
|
|
||||||
true);
|
|
||||||
done.run(CleanupBulkLoadResponse.newBuilder().build());
|
done.run(CleanupBulkLoadResponse.newBuilder().build());
|
||||||
} catch (IOException e) {
|
} catch (IOException e) {
|
||||||
ResponseConverter.setControllerException(controller, e);
|
ResponseConverter.setControllerException(controller, e);
|
||||||
|
|
|
@ -27,6 +27,7 @@ import java.io.IOException;
|
||||||
import java.util.TreeMap;
|
import java.util.TreeMap;
|
||||||
|
|
||||||
import org.apache.hadoop.conf.Configuration;
|
import org.apache.hadoop.conf.Configuration;
|
||||||
|
import org.apache.hadoop.fs.FileStatus;
|
||||||
import org.apache.hadoop.fs.FileSystem;
|
import org.apache.hadoop.fs.FileSystem;
|
||||||
import org.apache.hadoop.fs.Path;
|
import org.apache.hadoop.fs.Path;
|
||||||
import org.apache.hadoop.hbase.HBaseTestingUtility;
|
import org.apache.hadoop.hbase.HBaseTestingUtility;
|
||||||
|
@ -49,6 +50,7 @@ import org.junit.AfterClass;
|
||||||
import org.junit.BeforeClass;
|
import org.junit.BeforeClass;
|
||||||
import org.junit.Test;
|
import org.junit.Test;
|
||||||
import org.junit.experimental.categories.Category;
|
import org.junit.experimental.categories.Category;
|
||||||
|
import org.apache.hadoop.hbase.security.SecureBulkLoadUtil;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Test cases for the "load" half of the HFileOutputFormat bulk load
|
* Test cases for the "load" half of the HFileOutputFormat bulk load
|
||||||
|
@ -259,6 +261,16 @@ public class TestLoadIncrementalHFiles {
|
||||||
table.close();
|
table.close();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// verify staging folder has been cleaned up
|
||||||
|
Path stagingBasePath = SecureBulkLoadUtil.getBaseStagingDir(util.getConfiguration());
|
||||||
|
if(fs.exists(stagingBasePath)) {
|
||||||
|
FileStatus[] files = fs.listStatus(stagingBasePath);
|
||||||
|
for(FileStatus file : files) {
|
||||||
|
assertTrue("Folder=" + file.getPath() + " is not cleaned up.",
|
||||||
|
file.getPath().getName() != "DONOTERASE");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
util.deleteTable(tableName);
|
util.deleteTable(tableName);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
Loading…
Reference in New Issue