Compare commits
9 Commits
| Author | SHA1 | Date |
|---|---|---|
|
|
00259426be | |
|
|
ec691d979d | |
|
|
e573c43b7c | |
|
|
459aea75b1 | |
|
|
0e91b331c8 | |
|
|
28f5d8b5cd | |
|
|
c8c85cbcb0 | |
|
|
685af27030 | |
|
|
895fa84778 |
|
|
@ -24,5 +24,5 @@ s3.ak=ak
|
|||
s3.sk=sk
|
||||
s3.endpoint=endpoint
|
||||
s3.bucket_name=bucket
|
||||
s3.blocksize=1048576
|
||||
s3.chunksize=4194304
|
||||
s3.blocksize=4194304
|
||||
s3.chunksize=67108864
|
||||
|
|
|
|||
|
|
@ -38,7 +38,7 @@ enum TopoStatusCode {
|
|||
TOPO_IP_PORT_DUPLICATED = 14;
|
||||
TOPO_NAME_DUPLICATED = 15;
|
||||
TOPO_CREATE_COPYSET_ON_METASERVER_FAIL = 16;
|
||||
TOPO_CANNOT_REMOVE_NOT_RETIRED = 17;
|
||||
TOPO_CANNOT_REMOVE_NOT_OFFLINE = 17;
|
||||
TOPO_POOL_EXIST = 18;
|
||||
TOPO_LEADER_NOT_FOUND = 19;
|
||||
TOPO_PARTITION_NOT_FOUND = 20;
|
||||
|
|
|
|||
|
|
@ -151,15 +151,8 @@ CURVEFS_ERROR FuseClient::FuseOpInit(void *userdata,
|
|||
FSStatusCode ret = mdsClient_->GetFsInfo(fsName, &fsInfo);
|
||||
if (ret != FSStatusCode::OK) {
|
||||
if (FSStatusCode::NOT_FOUND == ret) {
|
||||
LOG(INFO) << "The fsName not exist, try to CreateFs"
|
||||
<< ", fsName = " << fsName;
|
||||
|
||||
CURVEFS_ERROR ret2 = CreateFs(userdata, &fsInfo);
|
||||
if (ret2 != CURVEFS_ERROR::OK) {
|
||||
LOG(ERROR) << "CreateFs failed, ret = " << ret2
|
||||
<< ", fsName = " << fsName;
|
||||
return ret2;
|
||||
}
|
||||
LOG(ERROR) << "The fsName not exist, fsName = " << fsName;
|
||||
return CURVEFS_ERROR::NOTEXIST;
|
||||
} else {
|
||||
LOG(ERROR) << "GetFsInfo failed, FSStatusCode = " << ret
|
||||
<< ", FSStatusCode_Name = "
|
||||
|
|
|
|||
|
|
@ -248,12 +248,18 @@ bool DiskCacheManager::IsDiskCacheFull() {
|
|||
}
|
||||
|
||||
bool DiskCacheManager::IsDiskCacheSafe() {
|
||||
int64_t ratio = diskFsUsedRatio_.load(std::memory_order_seq_cst);
|
||||
uint64_t usedBytes = GetDiskUsedbytes();
|
||||
if (usedBytes <= (safeRatio_ * maxUsableSpaceBytes_ / 100)) {
|
||||
if ((usedBytes < (safeRatio_ * maxUsableSpaceBytes_ / 100))
|
||||
&& (ratio < safeRatio_)) {
|
||||
VLOG(3) << "disk cache is safe"
|
||||
<< ", usedBytes is: " << usedBytes;
|
||||
<< ", usedBytes is: " << usedBytes
|
||||
<< ", use ratio is: " << ratio;
|
||||
return true;
|
||||
}
|
||||
VLOG(3) << "disk cache is not safe"
|
||||
<< ", usedBytes is: " << usedBytes
|
||||
<< ", use ratio is: " << ratio;
|
||||
return false;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -222,6 +222,8 @@ TopoStatusCode TopologyImpl::RemoveMetaServer(MetaServerIdType id) {
|
|||
WriteLockGuard wlockMetaServer(metaServerMutex_);
|
||||
auto it = metaServerMap_.find(id);
|
||||
if (it != metaServerMap_.end()) {
|
||||
uint64_t metaserverCapacity =
|
||||
it->second.GetMetaServerSpace().GetDiskCapacity();
|
||||
if (!storage_->DeleteMetaServer(id)) {
|
||||
return TopoStatusCode::TOPO_STORGE_FAIL;
|
||||
}
|
||||
|
|
@ -230,6 +232,17 @@ TopoStatusCode TopologyImpl::RemoveMetaServer(MetaServerIdType id) {
|
|||
ix->second.RemoveMetaServer(id);
|
||||
}
|
||||
metaServerMap_.erase(it);
|
||||
|
||||
// update pool
|
||||
WriteLockGuard wlockPool(poolMutex_);
|
||||
PoolIdType poolId = ix->second.GetPoolId();
|
||||
auto it = poolMap_.find(poolId);
|
||||
if (it != poolMap_.end()) {
|
||||
it->second.SetDiskCapacity(it->second.GetDiskCapacity() -
|
||||
metaserverCapacity);
|
||||
} else {
|
||||
return TopoStatusCode::TOPO_POOL_NOT_FOUND;
|
||||
}
|
||||
return TopoStatusCode::TOPO_OK;
|
||||
} else {
|
||||
return TopoStatusCode::TOPO_METASERVER_NOT_FOUND;
|
||||
|
|
|
|||
|
|
@ -317,6 +317,12 @@ void TopologyManager::DeleteServer(const DeleteServerRequest *request,
|
|||
<< ", serverId = " << request->serverid();
|
||||
response->set_statuscode(TopoStatusCode::TOPO_INTERNAL_ERROR);
|
||||
return;
|
||||
} else if (OnlineState::OFFLINE != ms.GetOnlineState()) {
|
||||
LOG(ERROR) << "Can not delete server which have "
|
||||
<< "metaserver not offline.";
|
||||
response->set_statuscode(
|
||||
TopoStatusCode::TOPO_CANNOT_REMOVE_NOT_OFFLINE);
|
||||
return;
|
||||
} else {
|
||||
errcode = topology_->RemoveMetaServer(msId);
|
||||
if (errcode != TopoStatusCode::TOPO_OK) {
|
||||
|
|
|
|||
|
|
@ -99,7 +99,7 @@ int CreateFsTool::Init() {
|
|||
mds::CreateFsRequest request;
|
||||
request.set_fsname(FLAGS_fsName);
|
||||
request.set_blocksize(FLAGS_blockSize);
|
||||
if (FLAGS_fsType == "s3") {
|
||||
if (FLAGS_fsType == kFsTypeS3) {
|
||||
// s3
|
||||
request.set_fstype(common::FSType::TYPE_S3);
|
||||
auto s3 = new common::S3Info();
|
||||
|
|
@ -110,7 +110,7 @@ int CreateFsTool::Init() {
|
|||
s3->set_blocksize(FLAGS_s3_blocksize);
|
||||
s3->set_chunksize(FLAGS_s3_chunksize);
|
||||
request.mutable_fsdetail()->set_allocated_s3info(s3);
|
||||
} else if (FLAGS_fsType == "volume") {
|
||||
} else if (FLAGS_fsType == kFsTypeVolume) {
|
||||
// volume
|
||||
request.set_fstype(common::FSType::TYPE_VOLUME);
|
||||
auto volume = new common::Volume();
|
||||
|
|
@ -121,14 +121,22 @@ int CreateFsTool::Init() {
|
|||
volume->set_password(FLAGS_volumePassword);
|
||||
request.mutable_fsdetail()->set_allocated_volume(volume);
|
||||
} else {
|
||||
std::cerr << "-fsType should be s3 or volume." << std::endl;
|
||||
std::cerr << "-fsType should be " << kFsTypeS3 << " or "
|
||||
<< kFsTypeVolume << "." << std::endl;
|
||||
ret = -1;
|
||||
}
|
||||
|
||||
AddRequest(request);
|
||||
|
||||
SetController();
|
||||
|
||||
return ret;
|
||||
}
|
||||
|
||||
void CreateFsTool::SetController() {
|
||||
controller_->set_timeout_ms(FLAGS_rpcTimeoutMs);
|
||||
}
|
||||
|
||||
bool CreateFsTool::AfterSendRequestToHost(const std::string& host) {
|
||||
bool ret = true;
|
||||
if (controller_->Failed()) {
|
||||
|
|
|
|||
|
|
@ -50,6 +50,7 @@ class CreateFsTool : public CurvefsToolRpc<curvefs::mds::CreateFsRequest,
|
|||
protected:
|
||||
void AddUpdateFlags() override;
|
||||
bool AfterSendRequestToHost(const std::string& host) override;
|
||||
void SetController() override;
|
||||
};
|
||||
|
||||
} // namespace create
|
||||
|
|
|
|||
|
|
@ -124,6 +124,21 @@ int CurvefsBuildTopologyTool::HandleBuildCluster() {
|
|||
return DealFailedRet(ret, "scan cluster");
|
||||
}
|
||||
|
||||
ret = RemoveServersNotInNewTopo();
|
||||
if (ret != 0) {
|
||||
return DealFailedRet(ret, "remove server");
|
||||
}
|
||||
|
||||
ret = RemoveZonesNotInNewTopo();
|
||||
if (ret != 0) {
|
||||
return DealFailedRet(ret, "remove zone");
|
||||
}
|
||||
|
||||
ret = RemovePoolsNotInNewTopo();
|
||||
if (ret != 0) {
|
||||
return DealFailedRet(ret, "remove pool");
|
||||
}
|
||||
|
||||
ret = CreatePool();
|
||||
if (ret != 0) {
|
||||
return DealFailedRet(ret, "create pool");
|
||||
|
|
@ -270,6 +285,8 @@ int CurvefsBuildTopologyTool::ScanCluster() {
|
|||
[it](Pool& data) { return data.name == it->poolname(); });
|
||||
if (ix != poolDatas.end()) {
|
||||
poolDatas.erase(ix);
|
||||
} else {
|
||||
poolToDel.emplace_back(it->poolid());
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -291,6 +308,8 @@ int CurvefsBuildTopologyTool::ScanCluster() {
|
|||
});
|
||||
if (ix != zoneDatas.end()) {
|
||||
zoneDatas.erase(ix);
|
||||
} else {
|
||||
zoneToDel.emplace_back(it->zoneid());
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -313,6 +332,8 @@ int CurvefsBuildTopologyTool::ScanCluster() {
|
|||
});
|
||||
if (ix != serverDatas.end()) {
|
||||
serverDatas.erase(ix);
|
||||
} else {
|
||||
serverToDel.emplace_back(it->serverid());
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -325,7 +346,6 @@ int CurvefsBuildTopologyTool::ListPool(std::list<PoolInfo>* poolInfos) {
|
|||
ListPoolResponse response;
|
||||
brpc::Controller cntl;
|
||||
cntl.set_timeout_ms(FLAGS_rpcTimeoutMs);
|
||||
cntl.set_log_id(1);
|
||||
|
||||
LOG(INFO) << "ListPool send request: " << request.DebugString();
|
||||
stub.ListPool(&cntl, &request, &response, nullptr);
|
||||
|
|
@ -357,7 +377,6 @@ int CurvefsBuildTopologyTool::GetZonesInPool(PoolIdType poolid,
|
|||
|
||||
brpc::Controller cntl;
|
||||
cntl.set_timeout_ms(FLAGS_rpcTimeoutMs);
|
||||
cntl.set_log_id(1);
|
||||
|
||||
LOG(INFO) << "ListZoneInPool, send request: " << request.DebugString();
|
||||
|
||||
|
|
@ -390,7 +409,6 @@ int CurvefsBuildTopologyTool::GetServersInZone(
|
|||
request.set_zoneid(zoneid);
|
||||
brpc::Controller cntl;
|
||||
cntl.set_timeout_ms(FLAGS_rpcTimeoutMs);
|
||||
cntl.set_log_id(1);
|
||||
|
||||
LOG(INFO) << "ListZoneServer, send request: " << request.DebugString();
|
||||
|
||||
|
|
@ -415,6 +433,115 @@ int CurvefsBuildTopologyTool::GetServersInZone(
|
|||
return 0;
|
||||
}
|
||||
|
||||
int CurvefsBuildTopologyTool::RemovePoolsNotInNewTopo() {
|
||||
TopologyService_Stub stub(&channel_);
|
||||
for (auto it : poolToDel) {
|
||||
DeletePoolRequest request;
|
||||
DeletePoolResponse response;
|
||||
request.set_poolid(it);
|
||||
|
||||
brpc::Controller cntl;
|
||||
cntl.set_timeout_ms(FLAGS_rpcTimeoutMs);
|
||||
|
||||
LOG(INFO) << "ClearPool, send request: " << request.DebugString();
|
||||
|
||||
stub.DeletePool(&cntl, &request, &response, nullptr);
|
||||
|
||||
if (cntl.ErrorCode() == EHOSTDOWN ||
|
||||
cntl.ErrorCode() == brpc::ELOGOFF) {
|
||||
return kRetCodeRedirectMds;
|
||||
} else if (cntl.Failed()) {
|
||||
LOG(ERROR) << "ClearPool errcorde = " << response.statuscode()
|
||||
<< ", error content:" << cntl.ErrorText()
|
||||
<< " , poolId = " << it;
|
||||
return kRetCodeCommonErr;
|
||||
}
|
||||
|
||||
if (response.statuscode() != TopoStatusCode::TOPO_OK) {
|
||||
LOG(ERROR) << "ClearPool rpc response fail. "
|
||||
<< "Message is :" << response.DebugString()
|
||||
<< " , poolId =" << it;
|
||||
return response.statuscode();
|
||||
} else {
|
||||
LOG(INFO) << "Received ClearPool response success, "
|
||||
<< response.DebugString();
|
||||
}
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
int CurvefsBuildTopologyTool::RemoveZonesNotInNewTopo() {
|
||||
TopologyService_Stub stub(&channel_);
|
||||
for (auto it : zoneToDel) {
|
||||
DeleteZoneRequest request;
|
||||
DeleteZoneResponse response;
|
||||
request.set_zoneid(it);
|
||||
|
||||
brpc::Controller cntl;
|
||||
cntl.set_timeout_ms(FLAGS_rpcTimeoutMs);
|
||||
|
||||
LOG(INFO) << "ClearZone, send request: " << request.DebugString();
|
||||
|
||||
stub.DeleteZone(&cntl, &request, &response, nullptr);
|
||||
|
||||
if (cntl.ErrorCode() == EHOSTDOWN ||
|
||||
cntl.ErrorCode() == brpc::ELOGOFF) {
|
||||
return kRetCodeRedirectMds;
|
||||
} else if (cntl.Failed()) {
|
||||
LOG(ERROR) << "ClearZone, errcorde = " << response.statuscode()
|
||||
<< ", error content:" << cntl.ErrorText()
|
||||
<< " , zoneId = " << it;
|
||||
return kRetCodeCommonErr;
|
||||
}
|
||||
if (response.statuscode() != TopoStatusCode::TOPO_OK) {
|
||||
LOG(ERROR) << "ClearZone Rpc response fail. "
|
||||
<< "Message is :" << response.DebugString()
|
||||
<< " , zoneId = " << it;
|
||||
return response.statuscode();
|
||||
} else {
|
||||
LOG(INFO) << "Received ClearZone Rpc success, "
|
||||
<< response.DebugString();
|
||||
}
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
int CurvefsBuildTopologyTool::RemoveServersNotInNewTopo() {
|
||||
TopologyService_Stub stub(&channel_);
|
||||
for (auto it : serverToDel) {
|
||||
DeleteServerRequest request;
|
||||
DeleteServerResponse response;
|
||||
request.set_serverid(it);
|
||||
|
||||
brpc::Controller cntl;
|
||||
cntl.set_timeout_ms(FLAGS_rpcTimeoutMs);
|
||||
|
||||
LOG(INFO) << "ClearServer, send request: " << request.DebugString();
|
||||
|
||||
stub.DeleteServer(&cntl, &request, &response, nullptr);
|
||||
|
||||
if (cntl.ErrorCode() == EHOSTDOWN ||
|
||||
cntl.ErrorCode() == brpc::ELOGOFF) {
|
||||
return kRetCodeRedirectMds;
|
||||
} else if (cntl.Failed()) {
|
||||
LOG(ERROR) << "ClearServer, errcorde = " << response.statuscode()
|
||||
<< ", error content : " << cntl.ErrorText()
|
||||
<< " , serverId = " << it;
|
||||
return kRetCodeCommonErr;
|
||||
}
|
||||
if (response.statuscode() != TopoStatusCode::TOPO_OK) {
|
||||
LOG(ERROR) << "ClearServer Rpc response fail. "
|
||||
<< "Message is :" << response.DebugString()
|
||||
<< " , serverId = " << it;
|
||||
return response.statuscode();
|
||||
} else {
|
||||
LOG(INFO) << "Received ClearServer Rpc success, "
|
||||
<< response.DebugString();
|
||||
}
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
int CurvefsBuildTopologyTool::CreatePool() {
|
||||
TopologyService_Stub stub(&channel_);
|
||||
for (auto it : poolDatas) {
|
||||
|
|
@ -431,7 +558,6 @@ int CurvefsBuildTopologyTool::CreatePool() {
|
|||
|
||||
brpc::Controller cntl;
|
||||
cntl.set_timeout_ms(FLAGS_rpcTimeoutMs);
|
||||
cntl.set_log_id(1);
|
||||
|
||||
LOG(INFO) << "CreatePool, send request: " << request.DebugString();
|
||||
|
||||
|
|
@ -470,7 +596,6 @@ int CurvefsBuildTopologyTool::CreateZone() {
|
|||
|
||||
brpc::Controller cntl;
|
||||
cntl.set_timeout_ms(FLAGS_rpcTimeoutMs);
|
||||
cntl.set_log_id(1);
|
||||
|
||||
LOG(INFO) << "CreateZone, send request: " << request.DebugString();
|
||||
|
||||
|
|
@ -485,7 +610,7 @@ int CurvefsBuildTopologyTool::CreateZone() {
|
|||
<< " , zoneName = " << it.name;
|
||||
return kRetCodeCommonErr;
|
||||
}
|
||||
if (response.statuscode() != 0) {
|
||||
if (response.statuscode() != TopoStatusCode::TOPO_OK) {
|
||||
LOG(ERROR) << "CreateZone Rpc response fail. "
|
||||
<< "Message is :" << response.DebugString()
|
||||
<< " , zoneName = " << it.name;
|
||||
|
|
@ -513,7 +638,6 @@ int CurvefsBuildTopologyTool::CreateServer() {
|
|||
|
||||
brpc::Controller cntl;
|
||||
cntl.set_timeout_ms(FLAGS_rpcTimeoutMs);
|
||||
cntl.set_log_id(1);
|
||||
|
||||
LOG(INFO) << "CreateServer, send request: " << request.DebugString();
|
||||
|
||||
|
|
|
|||
|
|
@ -133,6 +133,9 @@ class CurvefsBuildTopologyTool : public curvefs::tools::CurvefsTool {
|
|||
int InitPoolData();
|
||||
int ScanCluster();
|
||||
int ScanPool();
|
||||
int RemovePoolsNotInNewTopo();
|
||||
int RemoveZonesNotInNewTopo();
|
||||
int RemoveServersNotInNewTopo();
|
||||
int CreatePool();
|
||||
int CreateZone();
|
||||
int CreateServer();
|
||||
|
|
@ -150,6 +153,10 @@ class CurvefsBuildTopologyTool : public curvefs::tools::CurvefsTool {
|
|||
std::list<Zone> zoneDatas;
|
||||
std::list<Pool> poolDatas;
|
||||
|
||||
std::list<ServerIdType> serverToDel;
|
||||
std::list<ZoneIdType> zoneToDel;
|
||||
std::list<PoolIdType> poolToDel;
|
||||
|
||||
std::vector<std::string> mdsAddressStr_;
|
||||
int mdsAddressIndex_;
|
||||
brpc::Channel channel_;
|
||||
|
|
|
|||
|
|
@ -211,6 +211,7 @@ class CurvefsToolRpc : public CurvefsTool {
|
|||
return true;
|
||||
}
|
||||
controller_->Reset();
|
||||
SetController();
|
||||
}
|
||||
// send request to all host failed
|
||||
return false;
|
||||
|
|
|
|||
|
|
@ -57,6 +57,10 @@ DEFINE_string(s3_bucket_name, "bucketname", "s3 bucket name");
|
|||
DEFINE_uint64(s3_blocksize, 1048576, "s3 block size");
|
||||
DEFINE_uint64(s3_chunksize, 4194304, "s3 chunk size");
|
||||
|
||||
// list-topology
|
||||
DEFINE_string(jsonPath, "/tmp/topology.json", "output json path");
|
||||
DEFINE_string(jsonType, "build", "output json type(build or tree)");
|
||||
|
||||
// topology
|
||||
DEFINE_string(mds_addr, "127.0.0.1:6700",
|
||||
"mds ip and port, separated by \",\""); // NOLINT
|
||||
|
|
@ -215,6 +219,10 @@ std::function<bool(google::CommandLineFlagInfo*)> CheckPartitionIdDefault =
|
|||
std::bind(&CheckFlagInfoDefault<fLS::clstring>, std::placeholders::_1,
|
||||
"partitionId");
|
||||
|
||||
std::function<bool(google::CommandLineFlagInfo*)> CheckJsonPathDefault =
|
||||
std::bind(&CheckFlagInfoDefault<fLS::clstring>, std::placeholders::_1,
|
||||
"jsonPath");
|
||||
|
||||
/* translate to string */
|
||||
|
||||
auto StrVec2Str(const std::vector<std::string>& strVec) -> std::string {
|
||||
|
|
|
|||
|
|
@ -133,6 +133,14 @@ const char kHostFollowerValue[] = "follower";
|
|||
const char kEtcdLeaderValue[] = "StateLeader";
|
||||
const char kEtcdFollowerValue[] = "StateFollower";
|
||||
|
||||
/* fs type */
|
||||
const char kFsTypeS3[] = "s3";
|
||||
const char kFsTypeVolume[] = "volume";
|
||||
|
||||
/* json type */
|
||||
const char kJsonTypeBuild[] = "build";
|
||||
const char kJsonTypeTree[] = "tree";
|
||||
|
||||
} // namespace tools
|
||||
} // namespace curvefs
|
||||
|
||||
|
|
@ -155,6 +163,23 @@ const char kZone[] = "zone";
|
|||
const char kReplicasNum[] = "replicasnum";
|
||||
const char kCopysetNum[] = "copysetnum";
|
||||
const char kZoneNum[] = "zonenum";
|
||||
const char kClusterId[] = "clusterid";
|
||||
const char kPoolId[] = "poolid";
|
||||
const char kPoolName[] = "poolname";
|
||||
const char kCreateTime[] = "createtime";
|
||||
const char kPolicy[] = "policy";
|
||||
const char kPoollist[] = "poollist";
|
||||
const char kZonelist[] = "zonelist";
|
||||
const char kZoneId[] = "zoneid";
|
||||
const char kZoneName[] = "zonename";
|
||||
const char kServerlist[] = "serverlist";
|
||||
const char kServerId[] = "serverid";
|
||||
const char kHostName[] = "hostname";
|
||||
const char kMetaserverList[] = "metaserverlist";
|
||||
const char kMetaserverId[] = "metaserverid";
|
||||
const char kHostIp[] = "hostip";
|
||||
const char kPort[] = "port";
|
||||
const char kOnlineState[] = "state";
|
||||
|
||||
} // namespace topology
|
||||
} // namespace mds
|
||||
|
|
@ -245,6 +270,8 @@ extern std::function<bool(google::CommandLineFlagInfo*)> CheckCopysetIdDefault;
|
|||
extern std::function<bool(google::CommandLineFlagInfo*)>
|
||||
CheckPartitionIdDefault;
|
||||
|
||||
extern std::function<bool(google::CommandLineFlagInfo*)> CheckJsonPathDefault;
|
||||
|
||||
/* translate to string */
|
||||
std::string StrVec2Str(const std::vector<std::string>&);
|
||||
|
||||
|
|
|
|||
|
|
@ -23,9 +23,14 @@
|
|||
|
||||
#include <json/json.h>
|
||||
|
||||
#include <fstream>
|
||||
#include <memory>
|
||||
|
||||
#include "src/common/string_util.h"
|
||||
|
||||
DECLARE_string(mdsAddr);
|
||||
DECLARE_string(jsonPath);
|
||||
DECLARE_string(jsonType);
|
||||
|
||||
namespace curvefs {
|
||||
namespace tools {
|
||||
|
|
@ -33,8 +38,9 @@ namespace list {
|
|||
|
||||
void TopologyListTool::PrintHelp() {
|
||||
CurvefsToolRpc::PrintHelp();
|
||||
std::cout << " [-mdsAddr=" << FLAGS_mdsAddr << "]";
|
||||
std::cout << std::endl;
|
||||
std::cout << " [-mdsAddr=" << FLAGS_mdsAddr
|
||||
<< "] [-jsonType=" << FLAGS_jsonType
|
||||
<< " -jsonPath=" << FLAGS_jsonPath << "]" << std::endl;
|
||||
}
|
||||
|
||||
void TopologyListTool::AddUpdateFlags() {
|
||||
|
|
@ -69,6 +75,9 @@ bool TopologyListTool::AfterSendRequestToHost(const std::string& host) {
|
|||
ret = true;
|
||||
// clusterId
|
||||
clusterId_ = response_->clusterid();
|
||||
clusterId2CLusterInfo_.insert(std::pair<std::string, ClusterInfo>(
|
||||
clusterId_,
|
||||
ClusterInfo(clusterId_, std::vector<mds::topology::PoolIdType>())));
|
||||
// pool
|
||||
if (!GetPoolInfoFromResponse()) {
|
||||
ret = false;
|
||||
|
|
@ -89,29 +98,36 @@ bool TopologyListTool::AfterSendRequestToHost(const std::string& host) {
|
|||
// show
|
||||
if (show_) {
|
||||
// cluster
|
||||
std::cout << "[cluster]\nclusterId: " << clusterId_ << std::endl;
|
||||
std::cout << "[cluster]\n"
|
||||
<< mds::topology::kClusterId << ": " << clusterId_
|
||||
<< std::endl;
|
||||
|
||||
// pool
|
||||
std::cout << "[pool]" << std::endl;
|
||||
for (auto const& i : poolId2PoolInfo) {
|
||||
for (auto const& i : poolId2PoolInfo_) {
|
||||
ShowPoolInfo(i.second);
|
||||
}
|
||||
// zone
|
||||
std::cout << "[zone]" << std::endl;
|
||||
for (auto const& i : zoneId2ZoneInfo) {
|
||||
for (auto const& i : zoneId2ZoneInfo_) {
|
||||
ShowZoneInfo(i.second);
|
||||
}
|
||||
// server
|
||||
std::cout << "[server]" << std::endl;
|
||||
for (auto const& i : serverId2ServerInfo) {
|
||||
for (auto const& i : serverId2ServerInfo_) {
|
||||
ShowServerInfo(i.second);
|
||||
}
|
||||
// metaserver
|
||||
std::cout << "[metaserver]" << std::endl;
|
||||
for (auto const& i : metaserverId2MetaserverInfo) {
|
||||
for (auto const& i : metaserverId2MetaserverInfo_) {
|
||||
ShowMetaserverInfo(i.second);
|
||||
}
|
||||
}
|
||||
|
||||
google::CommandLineFlagInfo info;
|
||||
if (!CheckJsonPathDefault(&info) && ret) {
|
||||
OutputFile();
|
||||
}
|
||||
}
|
||||
return ret;
|
||||
}
|
||||
|
|
@ -119,9 +135,8 @@ bool TopologyListTool::AfterSendRequestToHost(const std::string& host) {
|
|||
PoolPolicy::PoolPolicy(const std::string& jsonStr) {
|
||||
Json::CharReaderBuilder reader;
|
||||
std::stringstream ss(jsonStr);
|
||||
std::string err;
|
||||
Json::Value json;
|
||||
bool parseCode = Json::parseFromStream(reader, ss, &json, &err);
|
||||
bool parseCode = Json::parseFromStream(reader, ss, &json, nullptr);
|
||||
if (parseCode && !json["replicaNum"].isNull() &&
|
||||
!json["copysetNum"].isNull() && !json["zoneNum"].isNull()) {
|
||||
replicaNum = json["replicaNum"].asUInt();
|
||||
|
|
@ -136,9 +151,9 @@ std::ostream& operator<<(std::ostream& os, const PoolPolicy& policy) {
|
|||
if (policy.error) {
|
||||
os << "policy has error!";
|
||||
} else {
|
||||
os << "copysetNum:" << policy.copysetNum
|
||||
<< " replicaNum:" << policy.replicaNum
|
||||
<< " zoneNum:" << policy.zoneNum;
|
||||
os << mds::topology::kCopysetNum << ":" << policy.copysetNum << " "
|
||||
<< mds::topology::kReplicasNum << ":" << policy.replicaNum << " "
|
||||
<< mds::topology::kZoneNum << ":" << policy.zoneNum;
|
||||
}
|
||||
return os;
|
||||
}
|
||||
|
|
@ -155,10 +170,11 @@ bool TopologyListTool::GetPoolInfoFromResponse() {
|
|||
ret = false;
|
||||
} else {
|
||||
for (auto const& i : pools.poolinfos()) {
|
||||
poolId2PoolInfo.insert(
|
||||
poolId2PoolInfo_.insert(
|
||||
std::pair<mds::topology::PoolIdType, PoolInfoType>(
|
||||
i.poolid(),
|
||||
PoolInfoType(i, std::vector<mds::topology::ZoneIdType>())));
|
||||
clusterId2CLusterInfo_[clusterId_].second.emplace_back(i.poolid());
|
||||
}
|
||||
}
|
||||
return ret;
|
||||
|
|
@ -176,13 +192,13 @@ bool TopologyListTool::GetZoneInfoFromResponse() {
|
|||
ret = false;
|
||||
} else {
|
||||
for (auto const& i : zones.zoneinfos()) {
|
||||
zoneId2ZoneInfo.insert(
|
||||
zoneId2ZoneInfo_.insert(
|
||||
std::pair<mds::topology::ZoneIdType, ZoneInfoType>(
|
||||
i.zoneid(),
|
||||
ZoneInfoType(i,
|
||||
std::vector<mds::topology::ServerIdType>())));
|
||||
if (poolId2PoolInfo.find(i.poolid()) != poolId2PoolInfo.end()) {
|
||||
poolId2PoolInfo[i.poolid()].second.emplace_back(i.zoneid());
|
||||
if (poolId2PoolInfo_.find(i.poolid()) != poolId2PoolInfo_.end()) {
|
||||
poolId2PoolInfo_[i.poolid()].second.emplace_back(i.zoneid());
|
||||
} else {
|
||||
errorOutput_ << "zone:" << i.zoneid()
|
||||
<< " has error: poolId:" << i.poolid()
|
||||
|
|
@ -206,14 +222,14 @@ bool TopologyListTool::GetServerInfoFromResponse() {
|
|||
ret = false;
|
||||
} else {
|
||||
for (auto const& i : servers.serverinfos()) {
|
||||
serverId2ServerInfo.insert(
|
||||
serverId2ServerInfo_.insert(
|
||||
std::pair<mds::topology::ServerIdType, ServerInfoType>(
|
||||
i.serverid(),
|
||||
ServerInfoType(
|
||||
i, std::vector<mds::topology::MetaServerIdType>())));
|
||||
|
||||
if (zoneId2ZoneInfo.find(i.zoneid()) != zoneId2ZoneInfo.end()) {
|
||||
zoneId2ZoneInfo[i.zoneid()].second.emplace_back(i.serverid());
|
||||
if (zoneId2ZoneInfo_.find(i.zoneid()) != zoneId2ZoneInfo_.end()) {
|
||||
zoneId2ZoneInfo_[i.zoneid()].second.emplace_back(i.serverid());
|
||||
} else {
|
||||
errorOutput_ << "server:" << i.serverid()
|
||||
<< " has error: zoneId:" << i.zoneid()
|
||||
|
|
@ -237,12 +253,12 @@ bool TopologyListTool::GetMetaserverInfoFromResponse() {
|
|||
ret = false;
|
||||
} else {
|
||||
for (auto const& i : metaservers.metaserverinfos()) {
|
||||
metaserverId2MetaserverInfo.insert(
|
||||
metaserverId2MetaserverInfo_.insert(
|
||||
std::pair<mds::topology::MetaServerIdType, MetaserverInfoType>(
|
||||
i.metaserverid(), MetaserverInfoType(i)));
|
||||
if (serverId2ServerInfo.find(i.serverid()) !=
|
||||
serverId2ServerInfo.end()) {
|
||||
serverId2ServerInfo[i.serverid()].second.emplace_back(
|
||||
if (serverId2ServerInfo_.find(i.serverid()) !=
|
||||
serverId2ServerInfo_.end()) {
|
||||
serverId2ServerInfo_[i.serverid()].second.emplace_back(
|
||||
i.metaserverid());
|
||||
} else {
|
||||
errorOutput_ << "metaserver:" << i.metaserverid()
|
||||
|
|
@ -290,6 +306,27 @@ void TopologyListTool::ShowMetaserverInfo(
|
|||
std::cout << MetaserverInfo2Str(metaserver) << std::endl;
|
||||
}
|
||||
|
||||
bool TopologyListTool::OutputFile() {
|
||||
Json::Value value;
|
||||
topology::TopologyTreeJson treeJson(*this);
|
||||
if (!treeJson.BuildJsonValue(&value, FLAGS_jsonType)) {
|
||||
std::cerr << "build json file failed!" << std::endl;
|
||||
return false;
|
||||
}
|
||||
|
||||
std::ofstream jsonFile;
|
||||
jsonFile.open(FLAGS_jsonPath.c_str(), std::ios::out);
|
||||
if (!jsonFile) {
|
||||
std::cerr << "open json file failed!" << std::endl;
|
||||
return false;
|
||||
}
|
||||
Json::StreamWriterBuilder clusterMap;
|
||||
std::unique_ptr<Json::StreamWriter> writer(clusterMap.newStreamWriter());
|
||||
writer->write(value, &jsonFile);
|
||||
jsonFile.close();
|
||||
return true;
|
||||
}
|
||||
|
||||
} // namespace list
|
||||
} // namespace tools
|
||||
} // namespace curvefs
|
||||
|
|
|
|||
|
|
@ -23,6 +23,7 @@
|
|||
|
||||
#include <brpc/channel.h>
|
||||
#include <gflags/gflags.h>
|
||||
#include <json/json.h>
|
||||
|
||||
#include <map>
|
||||
#include <string>
|
||||
|
|
@ -33,11 +34,19 @@
|
|||
#include "curvefs/src/mds/common/mds_define.h"
|
||||
#include "curvefs/src/tools/curvefs_tool.h"
|
||||
#include "curvefs/src/tools/curvefs_tool_define.h"
|
||||
#include "curvefs/src/tools/list/curvefs_topology_tree_json.h"
|
||||
|
||||
namespace curvefs {
|
||||
namespace tools {
|
||||
|
||||
namespace topology {
|
||||
class TopologyTreeJson;
|
||||
}
|
||||
|
||||
namespace list {
|
||||
|
||||
using ClusterInfo =
|
||||
std::pair<std::string, std::vector<mds::topology::PoolIdType>>;
|
||||
using PoolInfoType =
|
||||
std::pair<mds::topology::PoolInfo, std::vector<mds::topology::ZoneIdType>>;
|
||||
using ZoneInfoType = std::pair<mds::topology::ZoneInfo,
|
||||
|
|
@ -63,7 +72,6 @@ struct PoolPolicy {
|
|||
friend std::ostream& operator<<(std::ostream& os, const PoolPolicy& policy);
|
||||
};
|
||||
|
||||
// TODO(chengyi01): output a json file which can be used to build topology
|
||||
class TopologyListTool
|
||||
: public CurvefsToolRpc<curvefs::mds::topology::ListTopologyRequest,
|
||||
curvefs::mds::topology::ListTopologyResponse,
|
||||
|
|
@ -77,13 +85,17 @@ class TopologyListTool
|
|||
void PrintHelp() override;
|
||||
int Init() override;
|
||||
|
||||
friend class topology::TopologyTreeJson;
|
||||
|
||||
protected:
|
||||
bool OutputFile();
|
||||
void AddUpdateFlags() override;
|
||||
bool AfterSendRequestToHost(const std::string& host) override;
|
||||
|
||||
/**
|
||||
* @brief Get the PoolInfo From Response, fill into poolId2PoolInfo
|
||||
* not include zoneId list (will be filled in GetZoneInfoFromResponse)
|
||||
* will fill clusterId2CLusterInfo's poolId list
|
||||
*
|
||||
* @return true
|
||||
* @return false
|
||||
|
|
@ -135,26 +147,31 @@ class TopologyListTool
|
|||
|
||||
protected:
|
||||
std::string clusterId_;
|
||||
/**
|
||||
* poolId to clusterInfo and poolIds which belongs to pool
|
||||
*/
|
||||
std::map<std::string, ClusterInfo> clusterId2CLusterInfo_;
|
||||
|
||||
/**
|
||||
* @brief poolId to poolInfo and zoneIds which belongs to pool
|
||||
*
|
||||
* @details
|
||||
*/
|
||||
std::map<mds::topology::PoolIdType, PoolInfoType> poolId2PoolInfo;
|
||||
std::map<mds::topology::PoolIdType, PoolInfoType> poolId2PoolInfo_;
|
||||
|
||||
/**
|
||||
* @brief zoneId to zoneInfo and serverIds which belongs to zone
|
||||
*
|
||||
* @details
|
||||
*/
|
||||
std::map<mds::topology::ZoneIdType, ZoneInfoType> zoneId2ZoneInfo;
|
||||
std::map<mds::topology::ZoneIdType, ZoneInfoType> zoneId2ZoneInfo_;
|
||||
|
||||
/**
|
||||
* @brief serverId to serverInfo and metaserverIds which belongs to server
|
||||
*
|
||||
* @details
|
||||
*/
|
||||
std::map<mds::topology::ServerIdType, ServerInfoType> serverId2ServerInfo;
|
||||
std::map<mds::topology::ServerIdType, ServerInfoType> serverId2ServerInfo_;
|
||||
|
||||
/**
|
||||
* @brief metaserverId to metaserverInfo
|
||||
|
|
@ -162,7 +179,7 @@ class TopologyListTool
|
|||
* @details
|
||||
*/
|
||||
std::map<mds::topology::MetaServerIdType, MetaserverInfoType>
|
||||
metaserverId2MetaserverInfo;
|
||||
metaserverId2MetaserverInfo_;
|
||||
};
|
||||
|
||||
} // namespace list
|
||||
|
|
|
|||
|
|
@ -0,0 +1,208 @@
|
|||
/*
|
||||
* Copyright (c) 2022 NetEase Inc.
|
||||
*
|
||||
* Licensed 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.
|
||||
*/
|
||||
/*
|
||||
* Project: curve
|
||||
* Created Date: 2022-01-14
|
||||
* Author: chengyi01
|
||||
*/
|
||||
|
||||
#include "curvefs/src/tools/list/curvefs_topology_tree_json.h"
|
||||
|
||||
#include "curvefs/src/tools/list/curvefs_topology_list.h"
|
||||
|
||||
namespace curvefs {
|
||||
namespace tools {
|
||||
namespace topology {
|
||||
|
||||
TopologyTreeJson::TopologyTreeJson(const list::TopologyListTool& topologyTool) {
|
||||
clusterId_ = topologyTool.clusterId_;
|
||||
clusterId2CLusterInfo_ = topologyTool.clusterId2CLusterInfo_;
|
||||
poolId2PoolInfo_ = topologyTool.poolId2PoolInfo_;
|
||||
zoneId2ZoneInfo_ = topologyTool.zoneId2ZoneInfo_;
|
||||
serverId2ServerInfo_ = topologyTool.serverId2ServerInfo_;
|
||||
metaserverId2MetaserverInfo_ = topologyTool.metaserverId2MetaserverInfo_;
|
||||
} // namespace topology
|
||||
|
||||
bool TopologyTreeJson::BuildClusterMapPools(Json::Value* pools) {
|
||||
bool ret = true;
|
||||
for (auto const& i : poolId2PoolInfo_) {
|
||||
Json::Value value;
|
||||
value[mds::topology::kName] = i.second.first.poolname();
|
||||
|
||||
auto policy =
|
||||
list::PoolPolicy(i.second.first.redundanceandplacementpolicy());
|
||||
if (policy.error) {
|
||||
std::cerr << "build pool error!" << std::endl;
|
||||
ret = false;
|
||||
break;
|
||||
}
|
||||
value[mds::topology::kReplicasNum] = policy.replicaNum;
|
||||
value[mds::topology::kCopysetNum] = policy.copysetNum;
|
||||
value[mds::topology::kZoneNum] = policy.zoneNum;
|
||||
pools->append(value);
|
||||
}
|
||||
return ret;
|
||||
}
|
||||
|
||||
bool TopologyTreeJson::BuildClusterMapServers(Json::Value* servers) {
|
||||
bool ret = true;
|
||||
for (auto const& i : serverId2ServerInfo_) {
|
||||
Json::Value value;
|
||||
auto server = i.second.first;
|
||||
value[mds::topology::kName] = server.hostname();
|
||||
value[mds::topology::kInternalIp] = server.internalip();
|
||||
value[mds::topology::kInternalPort] = server.internalport();
|
||||
value[mds::topology::kExternalIp] = server.externalip();
|
||||
value[mds::topology::kExternalPort] = server.externalport();
|
||||
auto zone = zoneId2ZoneInfo_[server.zoneid()].first;
|
||||
value[mds::topology::kZone] = zone.zonename();
|
||||
auto pool = poolId2PoolInfo_[zone.poolid()].first;
|
||||
value[mds::topology::kPool] = pool.poolname();
|
||||
servers->append(value);
|
||||
}
|
||||
return ret;
|
||||
}
|
||||
|
||||
bool TopologyTreeJson::BuildJsonValue(Json::Value* value,
|
||||
const std::string& jsonType) {
|
||||
if (jsonType == kJsonTypeBuild) {
|
||||
// build json
|
||||
Json::Value pools;
|
||||
Json::Value servers;
|
||||
if (!(BuildClusterMapPools(&pools) &&
|
||||
BuildClusterMapServers(&servers))) {
|
||||
return false;
|
||||
}
|
||||
(*value)[mds::topology::kServers] = servers;
|
||||
(*value)[mds::topology::kPools] = pools;
|
||||
} else if (jsonType == kJsonTypeTree) {
|
||||
if (!GetClusterTree(value, clusterId_)) {
|
||||
return false;
|
||||
}
|
||||
} else {
|
||||
std::cerr << "-jsonType should be " << kJsonTypeBuild << " or "
|
||||
<< kJsonTypeTree << "!" << std::endl;
|
||||
return false;
|
||||
}
|
||||
|
||||
return true;
|
||||
}
|
||||
bool TopologyTreeJson::GetClusterTree(Json::Value* cluster,
|
||||
const std::string& clusterId) {
|
||||
bool ret = true;
|
||||
auto clusterInfo = clusterId2CLusterInfo_[clusterId];
|
||||
(*cluster)[mds::topology::kClusterId] = clusterInfo.first;
|
||||
for (auto const& poolId : clusterInfo.second) {
|
||||
Json::Value pool;
|
||||
ret = GetPoolTree(&pool, poolId);
|
||||
if (!ret) {
|
||||
break;
|
||||
}
|
||||
(*cluster)[mds::topology::kPoollist].append(pool);
|
||||
}
|
||||
return ret;
|
||||
}
|
||||
|
||||
bool TopologyTreeJson::GetPoolTree(Json::Value* pool, uint64_t poolId) {
|
||||
bool ret = true;
|
||||
auto poolInfo = poolId2PoolInfo_[poolId];
|
||||
(*pool)[mds::topology::kPoolId] = poolInfo.first.poolid();
|
||||
(*pool)[mds::topology::kPoolName] = poolInfo.first.poolname();
|
||||
(*pool)[mds::topology::kCreateTime] = poolInfo.first.createtime();
|
||||
Json::CharReaderBuilder reader;
|
||||
std::stringstream ss(poolInfo.first.redundanceandplacementpolicy());
|
||||
Json::Value policyValue;
|
||||
std::string err;
|
||||
if (!Json::parseFromStream(reader, ss, &policyValue, &err)) {
|
||||
std::cerr << "parse policy failed! error is " << err << std::endl;
|
||||
ret = false;
|
||||
}
|
||||
(*pool)[mds::topology::kPolicy] = policyValue;
|
||||
|
||||
for (auto const& zoneId : poolInfo.second) {
|
||||
Json::Value zone;
|
||||
ret = GetZoneTree(&zone, zoneId);
|
||||
if (!ret) {
|
||||
break;
|
||||
}
|
||||
(*pool)[mds::topology::kZonelist].append(zone);
|
||||
}
|
||||
|
||||
return ret;
|
||||
}
|
||||
|
||||
bool TopologyTreeJson::GetZoneTree(Json::Value* zone, uint64_t zoneId) {
|
||||
bool ret = true;
|
||||
auto zoneInfo = zoneId2ZoneInfo_[zoneId];
|
||||
(*zone)[mds::topology::kZoneId] = zoneInfo.first.zoneid();
|
||||
(*zone)[mds::topology::kZoneName] = zoneInfo.first.zonename();
|
||||
(*zone)[mds::topology::kPoolId] = zoneInfo.first.poolid();
|
||||
|
||||
for (auto const& serverId : zoneInfo.second) {
|
||||
Json::Value server;
|
||||
ret = GetServerTree(&server, serverId);
|
||||
if (!ret) {
|
||||
break;
|
||||
}
|
||||
(*zone)[mds::topology::kServerlist].append(server);
|
||||
}
|
||||
|
||||
return ret;
|
||||
}
|
||||
|
||||
bool TopologyTreeJson::GetServerTree(Json::Value* server, uint64_t serverId) {
|
||||
bool ret = true;
|
||||
auto serverinfo = serverId2ServerInfo_[serverId];
|
||||
(*server)[mds::topology::kServerId] = serverinfo.first.serverid();
|
||||
(*server)[mds::topology::kHostName] = serverinfo.first.hostname();
|
||||
(*server)[mds::topology::kInternalIp] = serverinfo.first.internalip();
|
||||
(*server)[mds::topology::kInternalPort] = serverinfo.first.internalport();
|
||||
(*server)[mds::topology::kExternalIp] = serverinfo.first.externalip();
|
||||
(*server)[mds::topology::kExternalPort] = serverinfo.first.externalport();
|
||||
(*server)[mds::topology::kZoneId] = serverinfo.first.zoneid();
|
||||
(*server)[mds::topology::kPoolId] = serverinfo.first.poolid();
|
||||
|
||||
for (auto const& metaserverId : serverinfo.second) {
|
||||
Json::Value metaserver;
|
||||
ret = GetMetaserverTree(&metaserver, metaserverId);
|
||||
if (!ret) {
|
||||
break;
|
||||
}
|
||||
(*server)[mds::topology::kMetaserverList].append(metaserver);
|
||||
}
|
||||
|
||||
return ret;
|
||||
}
|
||||
|
||||
bool TopologyTreeJson::GetMetaserverTree(Json::Value* metaserver,
|
||||
uint64_t metaserverId) {
|
||||
bool ret = true;
|
||||
auto metaserverinfo = metaserverId2MetaserverInfo_[metaserverId];
|
||||
(*metaserver)[mds::topology::kMetaserverId] = metaserverinfo.metaserverid();
|
||||
(*metaserver)[mds::topology::kHostName] = metaserverinfo.hostname();
|
||||
(*metaserver)[mds::topology::kHostIp] = metaserverinfo.hostip();
|
||||
(*metaserver)[mds::topology::kPort] = metaserverinfo.port();
|
||||
(*metaserver)[mds::topology::kExternalIp] = metaserverinfo.externalip();
|
||||
(*metaserver)[mds::topology::kExternalPort] = metaserverinfo.externalport();
|
||||
(*metaserver)[mds::topology::kOnlineState] = metaserverinfo.onlinestate();
|
||||
(*metaserver)[mds::topology::kServerId] = metaserverinfo.serverid();
|
||||
|
||||
return ret;
|
||||
}
|
||||
|
||||
} // namespace topology
|
||||
} // namespace tools
|
||||
} // namespace curvefs
|
||||
|
|
@ -0,0 +1,118 @@
|
|||
/*
|
||||
* Copyright (c) 2022 NetEase Inc.
|
||||
*
|
||||
* Licensed 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.
|
||||
*/
|
||||
|
||||
/*
|
||||
* Project: curve
|
||||
* Created Date: 2022-01-14
|
||||
* Author: chengyi01
|
||||
*/
|
||||
|
||||
#ifndef CURVEFS_SRC_TOOLS_LIST_CURVEFS_TOPOLOGY_TREE_JSON_H_
|
||||
#define CURVEFS_SRC_TOOLS_LIST_CURVEFS_TOPOLOGY_TREE_JSON_H_
|
||||
|
||||
#include <json/json.h>
|
||||
|
||||
#include <map>
|
||||
#include <string>
|
||||
#include <utility>
|
||||
#include <vector>
|
||||
|
||||
#include "curvefs/proto/topology.pb.h"
|
||||
#include "curvefs/src/mds/common/mds_define.h"
|
||||
#include "curvefs/src/tools/list/curvefs_topology_list.h"
|
||||
|
||||
namespace curvefs {
|
||||
namespace tools {
|
||||
|
||||
namespace list {
|
||||
class TopologyListTool;
|
||||
}
|
||||
|
||||
namespace topology {
|
||||
|
||||
using ClusterInfo =
|
||||
std::pair<std::string, std::vector<mds::topology::PoolIdType>>;
|
||||
using PoolInfoType =
|
||||
std::pair<mds::topology::PoolInfo, std::vector<mds::topology::ZoneIdType>>;
|
||||
using ZoneInfoType = std::pair<mds::topology::ZoneInfo,
|
||||
std::vector<mds::topology::ServerIdType>>;
|
||||
using ServerInfoType = std::pair<mds::topology::ServerInfo,
|
||||
std::vector<mds::topology::MetaServerIdType>>;
|
||||
using MetaserverInfoType = mds::topology::MetaServerInfo;
|
||||
|
||||
class TopologyTreeJson {
|
||||
public:
|
||||
explicit TopologyTreeJson(const list::TopologyListTool& topologyTool);
|
||||
|
||||
protected:
|
||||
std::string clusterId_;
|
||||
/**
|
||||
* poolId to clusterInfo and poolIds which belongs to pool
|
||||
*/
|
||||
std::map<std::string, ClusterInfo> clusterId2CLusterInfo_;
|
||||
|
||||
/**
|
||||
* @brief poolId to poolInfo and zoneIds which belongs to pool
|
||||
*
|
||||
* @details
|
||||
*/
|
||||
std::map<mds::topology::PoolIdType, PoolInfoType> poolId2PoolInfo_;
|
||||
|
||||
/**
|
||||
* @brief zoneId to zoneInfo and serverIds which belongs to zone
|
||||
*
|
||||
* @details
|
||||
*/
|
||||
std::map<mds::topology::ZoneIdType, ZoneInfoType> zoneId2ZoneInfo_;
|
||||
|
||||
/**
|
||||
* @brief serverId to serverInfo and metaserverIds which belongs to server
|
||||
*
|
||||
* @details
|
||||
*/
|
||||
std::map<mds::topology::ServerIdType, ServerInfoType> serverId2ServerInfo_;
|
||||
|
||||
/**
|
||||
* @brief metaserverId to metaserverInfo
|
||||
*
|
||||
* @details
|
||||
*/
|
||||
std::map<mds::topology::MetaServerIdType, MetaserverInfoType>
|
||||
metaserverId2MetaserverInfo_;
|
||||
|
||||
public:
|
||||
bool BuildClusterMapPools(Json::Value* pools);
|
||||
|
||||
bool BuildClusterMapServers(Json::Value* servers);
|
||||
|
||||
bool BuildJsonValue(Json::Value* value, const std::string& jsonType);
|
||||
|
||||
bool GetClusterTree(Json::Value* cluster, const std::string& clusterId);
|
||||
|
||||
bool GetPoolTree(Json::Value* pool, uint64_t poolId);
|
||||
|
||||
bool GetZoneTree(Json::Value* zone, uint64_t zoneId);
|
||||
|
||||
bool GetServerTree(Json::Value* server, uint64_t serverId);
|
||||
|
||||
bool GetMetaserverTree(Json::Value* metaserver, uint64_t metaserverId);
|
||||
};
|
||||
|
||||
} // namespace topology
|
||||
} // namespace tools
|
||||
} // namespace curvefs
|
||||
|
||||
#endif // CURVEFS_SRC_TOOLS_LIST_CURVEFS_TOPOLOGY_TREE_JSON_H_
|
||||
|
|
@ -292,7 +292,25 @@ TEST_F(TestDiskCacheManager, IsDiskCacheFull) {
|
|||
}
|
||||
|
||||
TEST_F(TestDiskCacheManager, IsDiskCacheSafe) {
|
||||
int ret = diskCacheManager_->IsDiskCacheSafe();
|
||||
S3ClientAdaptorOption option;
|
||||
option.diskCacheOpt.diskCacheType = (DiskCacheType)2;
|
||||
option.diskCacheOpt.cacheDir = "/mnt/test_unit";
|
||||
option.diskCacheOpt.trimCheckIntervalSec = 1;
|
||||
option.diskCacheOpt.fullRatio = 0;
|
||||
option.diskCacheOpt.safeRatio = 0;
|
||||
option.diskCacheOpt.maxUsableSpaceBytes = 0;
|
||||
option.diskCacheOpt.cmdTimeoutSec = 5;
|
||||
option.diskCacheOpt.asyncLoadPeriodMs = 10;
|
||||
S3Client *client = nullptr;
|
||||
diskCacheManager_->Init(client, option);
|
||||
bool ret = diskCacheManager_->IsDiskCacheSafe();
|
||||
ASSERT_EQ(false, ret);
|
||||
|
||||
option.diskCacheOpt.fullRatio = 100;
|
||||
option.diskCacheOpt.safeRatio = 99;
|
||||
option.diskCacheOpt.maxUsableSpaceBytes = 100000000;
|
||||
diskCacheManager_->Init(client, option);
|
||||
ret = diskCacheManager_->IsDiskCacheSafe();
|
||||
ASSERT_EQ(true, ret);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -188,29 +188,12 @@ TEST_F(TestFuseVolumeClient, FuseOpInit_when_fs_not_exist) {
|
|||
EXPECT_CALL(*mdsClient_, GetFsInfo(fsName, _))
|
||||
.WillOnce(Return(FSStatusCode::NOT_FOUND));
|
||||
|
||||
EXPECT_CALL(*blockDeviceClient_, Stat(volName, user, _))
|
||||
.WillOnce(Return(CURVEFS_ERROR::OK));
|
||||
|
||||
EXPECT_CALL(*mdsClient_, CreateFs(_, _, _))
|
||||
.WillOnce(Return(FSStatusCode::OK));
|
||||
|
||||
FsInfo fsInfoExp;
|
||||
fsInfoExp.set_fsid(100);
|
||||
fsInfoExp.set_fsname(fsName);
|
||||
EXPECT_CALL(*mdsClient_, MountFs(fsName, _, _))
|
||||
.WillOnce(DoAll(SetArgPointee<2>(fsInfoExp), Return(FSStatusCode::OK)));
|
||||
|
||||
EXPECT_CALL(*blockDeviceClient_, Open(volName, user))
|
||||
.WillOnce(Return(CURVEFS_ERROR::OK));
|
||||
|
||||
CURVEFS_ERROR ret = client_->FuseOpInit(&mOpts, nullptr);
|
||||
ASSERT_EQ(CURVEFS_ERROR::OK, ret);
|
||||
|
||||
auto fsInfo = client_->GetFsInfo();
|
||||
ASSERT_NE(fsInfo, nullptr);
|
||||
|
||||
ASSERT_EQ(fsInfo->fsid(), fsInfoExp.fsid());
|
||||
ASSERT_EQ(fsInfo->fsname(), fsInfoExp.fsname());
|
||||
ASSERT_EQ(CURVEFS_ERROR::NOTEXIST, ret);
|
||||
}
|
||||
|
||||
TEST_F(TestFuseVolumeClient, FuseOpDestroy) {
|
||||
|
|
@ -1768,23 +1751,12 @@ TEST_F(TestFuseS3Client, FuseOpInit_when_fs_not_exist) {
|
|||
EXPECT_CALL(*mdsClient_, GetFsInfo(fsName, _))
|
||||
.WillOnce(Return(FSStatusCode::NOT_FOUND));
|
||||
|
||||
EXPECT_CALL(*mdsClient_, CreateFsS3(_, _, _))
|
||||
.WillOnce(Return(FSStatusCode::OK));
|
||||
|
||||
FsInfo fsInfoExp;
|
||||
fsInfoExp.set_fsid(100);
|
||||
fsInfoExp.set_fsname(fsName);
|
||||
EXPECT_CALL(*mdsClient_, MountFs(fsName, _, _))
|
||||
.WillOnce(DoAll(SetArgPointee<2>(fsInfoExp), Return(FSStatusCode::OK)));
|
||||
|
||||
CURVEFS_ERROR ret = client_->FuseOpInit(&mOpts, nullptr);
|
||||
ASSERT_EQ(CURVEFS_ERROR::OK, ret);
|
||||
|
||||
auto fsInfo = client_->GetFsInfo();
|
||||
ASSERT_NE(fsInfo, nullptr);
|
||||
|
||||
ASSERT_EQ(fsInfo->fsid(), fsInfoExp.fsid());
|
||||
ASSERT_EQ(fsInfo->fsname(), fsInfoExp.fsname());
|
||||
ASSERT_EQ(CURVEFS_ERROR::NOTEXIST, ret);
|
||||
}
|
||||
|
||||
TEST_F(TestFuseS3Client, FuseOpDestroy) {
|
||||
|
|
|
|||
|
|
@ -777,6 +777,48 @@ TEST_F(TestTopologyManager, test_DeleteServer_success) {
|
|||
ASSERT_EQ(TopoStatusCode::TOPO_OK, response.statuscode());
|
||||
}
|
||||
|
||||
TEST_F(TestTopologyManager, test_DeleteServerHaveMetaserver_success) {
|
||||
PoolIdType poolId = 0x11;
|
||||
ZoneIdType zoneId = 0x21;
|
||||
ServerIdType serverId = 0x31;
|
||||
PrepareAddPool(poolId);
|
||||
PrepareAddZone(zoneId);
|
||||
PrepareAddServer(serverId, "hostname1", "ip1", 0, "ip2", 0, zoneId, poolId);
|
||||
PrepareAddMetaServer(0x41, "ms1", "token1", 0x31, "ip1", 0, "ip2", 8888,
|
||||
OnlineState::OFFLINE);
|
||||
|
||||
DeleteServerRequest request;
|
||||
request.set_serverid(serverId);
|
||||
|
||||
DeleteServerResponse response;
|
||||
|
||||
EXPECT_CALL(*storage_, DeleteMetaServer(_)).WillOnce(Return(true));
|
||||
EXPECT_CALL(*storage_, DeleteServer(_)).WillOnce(Return(true));
|
||||
|
||||
serviceManager_->DeleteServer(&request, &response);
|
||||
|
||||
ASSERT_EQ(TopoStatusCode::TOPO_OK, response.statuscode());
|
||||
}
|
||||
|
||||
TEST_F(TestTopologyManager, test_DeleteServerHaveMetaserver_fail) {
|
||||
PoolIdType poolId = 0x11;
|
||||
ZoneIdType zoneId = 0x21;
|
||||
ServerIdType serverId = 0x31;
|
||||
PrepareAddPool(poolId);
|
||||
PrepareAddZone(zoneId);
|
||||
PrepareAddServer(serverId, "hostname1", "ip1", 0, "ip2", 0, zoneId, poolId);
|
||||
PrepareAddMetaServer(0x41, "ms1", "token1", 0x31, "ip1", 0, "ip2", 8888);
|
||||
DeleteServerRequest request;
|
||||
request.set_serverid(serverId);
|
||||
|
||||
DeleteServerResponse response;
|
||||
|
||||
serviceManager_->DeleteServer(&request, &response);
|
||||
|
||||
ASSERT_EQ(TopoStatusCode::TOPO_CANNOT_REMOVE_NOT_OFFLINE,
|
||||
response.statuscode());
|
||||
}
|
||||
|
||||
TEST_F(TestTopologyManager, test_ListZoneServer_ByIdSuccess) {
|
||||
PoolIdType poolId = 0x11;
|
||||
ZoneIdType zoneId = 0x21;
|
||||
|
|
|
|||
|
|
@ -0,0 +1,107 @@
|
|||
# CurveFS MetaServer
|
||||
|
||||
## 概述
|
||||
|
||||
MetaServer 在 CurveFS 集群中提供高可用、高可靠的元数据服务,并保证文件系统元数据的一致性。同时,在设计之初就以高性能和扩展性作为目标。
|
||||
|
||||
MetaServer 在整体设计上参考了 CurveBS 的 ChunkServer,单个 MetaServer 以用户态进程的形式运行在宿主机上,在 CPU/RAM 等资源充足的情况下,一台宿主机可以运行多个 MetaServer 进程。同时,也引入了 ChunkServer 中 Copyset 的设计,利用 Raft 保证数据的一致性和服务的高可用。
|
||||
|
||||
在元数据管理层面,对文件系统元数据进行分片管理,避免单个 Raft Group 维护一个文件系统元数据是带来的性能瓶颈。元数据的每个分片称为 Partition。Copyset 与 Partition 的对应关系可以是一对一,也可以是一对多。一对多的情况下表示,一个 Copyset 可以维护多个 Partition。在一对多的情况下,文件系统的元数据管理如下图所示:
|
||||
|
||||

|
||||
|
||||
图中共有两个 Copyset,三个副本放置在三台机器上。P1/P2/P3/P4 表示文件系统的元数据分片,其中 P1/P3 属于一个文件系统,P2/P4 属于一个文件系统。
|
||||
|
||||
## 整体架构
|
||||
|
||||
MetaServer 的整体架构如下所示,大致可分为三个部分:Service Layer、Core Business Layer 和 MetaStore。三者之间相互协作,高效处理外部组件的各种请求和任务。下面将对各个模块进行详细的介绍。
|
||||
|
||||

|
||||
|
||||
### Service Layer
|
||||
|
||||
对外提供 RPC 接口,供系统中的其他服务(Curve-Fuse、MDS、MetaServer 等)调用。同时提供 RESTful 接口,可以把当前进程中各组件的状态(Metric)同步给 Prometheus。
|
||||
|
||||
#### MetaService
|
||||
|
||||
提供元数据服务,是 MetaServer 的核心服务,提供文件系统元数据查询、创建、更新、删除操作的必要接口,例如 CreateInode、CreateDentry、ListDentry,并支持 Partition 的动态创建和删除。
|
||||
|
||||
#### CopysetService
|
||||
|
||||
提供动态创建Copyset、查询Copyset状态接口。在创建文件系统时,MDS会根据当前集群的负载情况,决定是否创建新的Copyset。
|
||||
|
||||
#### RaftService
|
||||
|
||||
由 [braft](https://github.com/baidu/braft) 提供,用于 Raft 一致性协议的交互。
|
||||
|
||||
#### CliService (Command Line Service)
|
||||
|
||||
提供 Raft 配置变更接口,包括 AddPeer、RemovePeer、ChangePeer、TransferLeader。同时提供了一个额外的 GetLeader 接口,用于获取当前复制组最新的 Leader 信息。
|
||||
|
||||
#### MetricService
|
||||
|
||||
提供 RESTful 接口,可以获取进程各组件的状态,Prometheus 会调用该接口采集数据,并利用 Grafana 可视化展示。
|
||||
|
||||
### Core Business Layer
|
||||
|
||||
MetaServer 核心处理逻辑,包括元数据请求的处理,并保证元数据的一致性、高可用、高可靠;心跳上报及配置变更任务的执行处理;以及注册模块等。
|
||||
|
||||
#### CopysetNode
|
||||
|
||||
表示 Raft Group 中的一个副本,是 braft raft node 的简单封装,同时实现了 Raft 状态机。
|
||||
|
||||
#### ApplyQueue
|
||||
|
||||
用于隔离 braft apply 线程,可以 apply 的请求会放入 ApplyQueue 中,同时 ApplyQueue 保证请求的有序执行,在请求执行完成后,返回给客户端响应。
|
||||
|
||||
#### MetaOperator
|
||||
|
||||
元数据请求到达后,会生成一个对应的 operator,operator 会将请求封装成 task,然后交给元数据请求对应的 CopysetNode 进行处理,完成副本间的数据同步。
|
||||
|
||||
#### Register
|
||||
|
||||
正常集群的启动流程是,先启动MDS,然后创建逻辑池,最后启动 MetaServer。在创建逻辑池时,需要指定逻辑池的拓扑结构,以及各个 MetaServer 进程的 IP 和 Port。这样做的目的是,阻止非法的 MetaServer 加入集群。
|
||||
|
||||
所以,在 MetaServer 的启动阶段,需要先向 MDS 进行注册,MDS 验证后会返回唯一标识 MetaServerID 及 Token。在后续 MetaServer 与 MDS 通讯时,需要提供此 ID 和 Token 作为身份标识和合法性认证信息。
|
||||
|
||||
#### Heartbeat
|
||||
|
||||
MDS 需要实时的信息来确认 MetaServer 的在线状态,并获取 MetaServer 和 Copyset 的状态和统计数据,并根据所有 MetaServer 的信息计算当前集群是否需要动态调度以及相应的调度命令。
|
||||
|
||||
MetaServer 以心跳的方式来完成上述功能,通过周期性的心跳,上报 MetaServer 和 Copyset 的信息,同时执行心跳响应中的调度任务。
|
||||
|
||||
#### Metric
|
||||
|
||||
利用 [bvar](https://github.com/apache/incubator-brpc/blob/master/docs/en/bvar.md) 导出系统中核心模块的统计信息。
|
||||
|
||||
### MetaStore
|
||||
|
||||
高效组织和管理内存元数据,同时配合 Raft 对元数据进行定期 dump,加速重启过程。
|
||||
|
||||
#### MetaPartition
|
||||
|
||||
文件系统的元数据进行分片管理,每个分片称为 Partition,Partition 通过聚合 InodeStorage 和 DentryStorage,提供了对 Dentry 和 Inode 的增删改查接口,同时 Partition 管理的元数据全部缓存在内存中。
|
||||
|
||||
Inode 对应文件系统中的一个文件或目录,记录相应的元数据信息,比如 atime/ctime/mtime 等。当 Inode 表示一个文件时,还会记录文件的数据寻址信息。每个 Partition 管理固定范围内的 Inode,根据 InodeId 进行划分,比如 InodeId [1-200] 由 Partition 1管理,InodeId [201-400] 由 Partition 2 管理,依次类推。
|
||||
|
||||
Dentry 是文件系统中的目录项,记录文件名到 inode 的映射关系。一个父目录下所有文件/目录的 Dentry 信息由父目录 Inode 所在的 Partition 进行管理。
|
||||
|
||||
#### MetaSnapshot
|
||||
|
||||
配合 CopysetNode 实现 Raft 快照,定期将 MetaPartition 中记录的元数据信息 dump 到本地磁盘上,起到了启动加速以及元数据去重的功能。
|
||||
|
||||
当 Raft 快照触发时,MetaStore 会 fork 出一个子进程,子进程会把当前 MetaPartition 中记录的所有元数据进行序列化并持久化到本地磁盘上。在进程重启时,会首先加载上次的 Raft 快照到 MetaPartition 中,然后再从 Raft 日志中回放元数据操作记录。
|
||||
|
||||
#### S3Compaction
|
||||
|
||||
在对接 S3 文件系统(文件系统数据存放到 S3)时,由于大多数 S3 服务不支持对象的覆盖写/追加写,所以在对文件进行上述写入时,Curve-Fuse 会把新写入的数据上传到一个新的 S3 对象,并向 Inode 的 extent 字段中中插入一条相应的记录。
|
||||
|
||||
以下图为例,用户在第一次写入文件后,进行了三次覆盖写,所以 extent 中会记录 4 条记录。在没有 Compaction 的情况下,后续的读操作需要计算出每个范围的最新数据,然后分别从 S3 下载、合并,最终返回给上层应用。这里的性能开销、空间浪费是显而易见的,但是上层应用的写入模式是无法限制的。
|
||||
|
||||

|
||||
|
||||
Compaction 的主要作用就是将 extent 中有重叠或连续的写入进行合并,生成一个新的 S3 对象,以加快后续的读取速度,减少存储空间浪费。
|
||||
|
||||
#### Trash
|
||||
|
||||
在当前的设计中,Inode 的 nlink 计数减到 0 时,并没有立即对该 Inode 进行清理,而是将 Inode 标记为待清理状态,由 Trash 模块进行定期扫描,当超过预设的阈值时间后,才将 Inode 从 MetaPartition 中删除。
|
||||
|
|
@ -0,0 +1,251 @@
|
|||
# CurveFS 空间分配 POC 版本设计方案
|
||||
|
||||
## 背景
|
||||
|
||||
根据CurveFS总体方案设计,文件系统基于当前的块进行实现,所以需要设计基于块的空间分配器,用于分配并存储文件数据。
|
||||
|
||||
## 本地文件系统空间分配相关特性
|
||||
|
||||
### 局部性
|
||||
|
||||
尽量分配连续的磁盘空间,存储文件的数据。这一特性主要是针对 HDD 进行的优化,降低磁盘寻道时间。
|
||||
|
||||
### 延迟分配/Allocate-on-flush
|
||||
|
||||
在 sync/flush 之前,尽可能多的积累更多的文件数据块才进行空间分配,一方面可以提高局部性,另一方面可以降低磁盘碎片。
|
||||
|
||||
### Inline file/data
|
||||
|
||||
几百字节的小文件不单独分配磁盘空间,直接把数据存放到文件的元数据中。
|
||||
|
||||
针对上述的本地文件系统特性,Curve 文件系统分配需要着重考虑**局部性**。
|
||||
|
||||
虽然 Curve 是一个分布式文件系统,但是单个文件系统的容量可能会比较大,如果在空间分配时,不考虑局部性,inode 中记录的 extent 数量很多,导致文件系统元数据量很大。
|
||||
|
||||
假如文件系统大小为 1PiB,空间分配粒度为 1MiB,inode 中存储的 extent为三元组(fileoffset,blockoffset,length),当空间完全分配之后,extent 的元数据量为 24GiB(1PiB / 1MiB * 24,24 为每个 extent 所占用的字节大小)。
|
||||
|
||||
如果同一文件在多次申请空间时,能分配连续的地址空间,则 extent 可以进行合并。例如,文件先后写入两次,每次写入 1MiB 数据,分别申请的地址空间为(100MiB,1MiB)和(101MiB,1MiB),则只需要一个 extent 进行记录即可,(0,100MiB,2MiB)。
|
||||
|
||||
所以,如果能对文件的多次空间申请分配连续的地址空间,则 inode 中记录的 extent 数量可以大大减少,能够降低整个文件系统的元数据量。
|
||||
|
||||
对于延迟分配和 Inline file 这两个特性,需要 Curve Fuse 端配合完成。
|
||||
|
||||
## 空间分配
|
||||
|
||||
### 整体设计
|
||||
|
||||

|
||||
|
||||
分配器包括两层结构:
|
||||
|
||||
第一层用 bitmap 进行表示,每个 bit 标识其所对应的一块空间(以 4MiB 为例,具体大小可配置)是否分配出去。
|
||||
|
||||
第二层为 free extent list,表示每个已分配的块,哪些仍然是空闲的(offset, length),以 offset 为 key 进行排序(这里可以用 map 或者 btree 对所有的 free extent 进行管理)。
|
||||
|
||||
当前设计不考虑持久化问题,空间分配器只作为内存结构,负责空间的分配与回收。在初始化时,扫描文件系统所有 inode 中已使用的空间。
|
||||
|
||||
### 空间分配流程
|
||||
|
||||
在新文件进行空间分配时,随机选择 level1 中标记为 0 的块,**先预分配给这个文件,但是并不表示这个块被该文件独占**。
|
||||
|
||||
以下图为例:file1 新申请了 2MiB 的空间。首先从 level1 中随机选一个标记为 0 的块分配出去,然后将这一个块中的前 2MiB 空间分配给这个文件,剩余部分加入到 level2 中的 list 中。
|
||||
|
||||

|
||||
|
||||
后续,file1 再次追加写入 2MiB 数据,此时申请空间时,需要附带上 file1 最后一个字节数据在底层存储的位置,再加 1(期望申请的地址空间起始 offset)。以图中为例,则附带的值为 30MiB。
|
||||
|
||||
这次的空间申请,直接从 level2 中以 30MiB 作为 key 进行查找,找到后,进行空间分配。分配之后,相关信息如下图所示:
|
||||
|
||||

|
||||
|
||||
之前剩余的 30MiB ~ 2MiB 的 extent 完全分配出去,所以从 level2 中的 list 中删除。
|
||||
|
||||
文件 inode 中的 extent 可以将两次的申请结果进行合并,得到(0,28MiB,4MiB)。
|
||||
|
||||
#### 特殊情况
|
||||
|
||||
1. 新文件申请空间时,level1 中的所有 bit 都标记为 1,即所有的块都已经预分配出去。在文件系统空间比较满的情况下,有可能会造成这个问题。此时,申请空间时,需要从 level2 中,随机或者选择可用空间最大的 extent 分配出去。
|
||||
|
||||
2. 文件申请空间时,之前预分配块的剩余空间被其他文件占用。此时,首先从 level1 查找一个可用的块,不满足要求时,按情况 1 进行处理。
|
||||
|
||||
3. file1 再次追加写入数据时,会附带 32MiB 来申请空间。此时,从 level1 中查找 32MiB 对应的块标记是否为 0,如果为 0,则将这个块继续分配给 file1。否则,可以从 level1 中随机选择一个可用的块进行分配。尽可能合并多个块分配给同一个文件。
|
||||
|
||||
### 空间回收
|
||||
|
||||
空间回收主要是一个extent合并的过程,有以下几种情况:
|
||||
|
||||
1. 文件释放了一个完整的块,则直接将level1中对应的bit置为0。
|
||||
|
||||
2. 文件释放了一小段空间,则尝试与level2中的extent进行合并。
|
||||
|
||||
a. 如果合并之后是一个完整的块,则重新将level1中对应的bit置为0,同时删除该extent。
|
||||
|
||||
b. 如果不能合并,则向level2中插入一个新的extent。
|
||||
|
||||
|
||||
### 小文件处理
|
||||
|
||||
大量小文件的情况下,按照上述的分配策略,会导致 level1 的 bitmap 标记全为 1,同时 level2 中也会有很多 extent。
|
||||
|
||||
所以可以参考 [chubaofs](https://github.com/chubaofs/chubaofs),对大小文件区分不同的分配逻辑。同时,将文件系统的空间划分成两个部分,一部分用于小文件的空间分配,另一部分用于大文件分配。两部分空间是相对的,一部分用完后,可以申请另一部分的空间。比如,大文件部分的空间完全分配出去,则可以继续从小文件空间进行分配。
|
||||
|
||||
用于小文件空间分配的部分,空闲空间可以用 extent 来表示。
|
||||
|
||||

|
||||
|
||||
小文件在空间分配时,也需要考虑尽量分配连续的地址空间。
|
||||
|
||||
文件在第一次申请空间时,选择一个能满足要求的 extent 分配出去。后续的空间申请,同样要带上文件最后一个字节所在的地址空间,用于尽量分配连续的地址空间。
|
||||
|
||||
文件空间的申请,具体由大文件,还是由小文件处理,可以参考如下策略,大小文件阈值为 1MiB:
|
||||
|
||||

|
||||
|
||||
### 并发问题
|
||||
|
||||
如果所有的空间分配和回收全部由一个分配器来进行管理,那么这里的分配很有可能成为一个瓶颈。
|
||||
|
||||
为了避免整个问题,可以将整个空间,由多个分配器来进行管理,每个分配器管理不同的地址空间。比如,将整个空间划分为 10 组,每组空间都有一个空间分配器进行管理。
|
||||
|
||||
在申请空间时,如果没有附带期望地址空间的 offset,则随机选取一个分配器进行空间分配。如果附带了期望的 offset,则由对应的分配器进行处理。
|
||||
|
||||
空间回收时,根据回收的 offset,交给对应的分配器去回收。
|
||||
|
||||
### 文件系统扩容
|
||||
|
||||
在线扩容时,直接在新扩容的空间上,创建新的空间分配器进行空间管理。
|
||||
|
||||
文件系统重新加载时,再将所有的空间,按照上述的策略,进行分组管理。
|
||||
|
||||
### 接口设计
|
||||
|
||||
#### RPC 接口
|
||||
|
||||
当前设计是把空间分配器作为内置服务放在元数据节点,所以请求的发起方是 Curve Fuse,元数据服务器接收到请求后,根据 fsId 查找到对应的文件系统的空间分配器后,将空间分配/回收的任务交给这个分配器进行处理,处理完成后,返回 RPC。
|
||||
|
||||
空间分配器相关的 RPC 接口,及 request/response 定义如下。
|
||||
|
||||
```protobuf
|
||||
syntax="proto2";
|
||||
option cc_generic_services = true;
|
||||
|
||||
enum StatusCode {
|
||||
UNKNOWN_ERROR = 0; // 未知错误
|
||||
OK = 1; // 成功
|
||||
NOSPACE = 2; // 空间不足
|
||||
}
|
||||
message Extent {
|
||||
required uint64 offset = 1; // 块设备地址空间起始地址
|
||||
required uint32 length = 2; // 长度
|
||||
}
|
||||
enum AllocateType {
|
||||
NONE = 0;
|
||||
SMALL = 1; // 小文件分配
|
||||
BIG = 2; // 大文件分配
|
||||
}
|
||||
|
||||
message AllocateHint {
|
||||
optional AllocateType allocType = 1; // 申请类型
|
||||
optional uint64 leftOffset = 2; // 期望申请到的地址空间的起始位置
|
||||
optional uint64 rightOffset = 3; // 期望申请到的地址空间的结束位置
|
||||
}
|
||||
|
||||
message AllocateSpaceRequest {
|
||||
required uint64 fsId = 1; // 文件系统ID
|
||||
required uint32 size = 2; // 申请空间的大小
|
||||
optional AllocateHint allocHint = 3;
|
||||
}
|
||||
|
||||
message AllocateSpaceResponse {
|
||||
required StatusCode status = 1; // 状态码
|
||||
repeated Extent extents = 2; // 申请到的地址空间,可能不连续,所以以repeated表示
|
||||
}
|
||||
|
||||
message DeallocateSpaceRequest {
|
||||
required uint64 fsId = 1;
|
||||
repeated Extent extents = 2; // 释放的空间
|
||||
}
|
||||
|
||||
message DeallocateSpaceResponse {
|
||||
required StatusCode status = 1;
|
||||
}
|
||||
|
||||
service MetaServerService {
|
||||
rpc AllocateSpace(AllocateSpaceRequest) returns (AllocateSpaceResponse);
|
||||
rpc DeallocateSpace(DeallocateSpaceRequest) returns (DeallocateSpaceResponse);
|
||||
}
|
||||
```
|
||||
|
||||
#### 空间分配器接口
|
||||
|
||||
空间分配器相关接口及部分数据结构定义如下:
|
||||
|
||||
```cpp
|
||||
#include <cstdint>
|
||||
#include <vector>
|
||||
|
||||
enum class AllocateType {
|
||||
NONE = 0,
|
||||
SMALL = 1,
|
||||
BIG = 2
|
||||
};
|
||||
|
||||
struct AllocateHint {
|
||||
AllocateType allocType = AllocateType::NONE;
|
||||
uint64_t leftOffset = 0;
|
||||
uint64_t rightOffset = 0;
|
||||
};
|
||||
|
||||
struct Extent {
|
||||
uint64_t offset = 0;
|
||||
uint32_t len = 0;
|
||||
};
|
||||
|
||||
using Extents = std::vector<Extent>;
|
||||
|
||||
class Allocator {
|
||||
public:
|
||||
Allocator(...) {}
|
||||
virtual ~Allocator() = default;
|
||||
|
||||
/**
|
||||
* @brief 申请空间
|
||||
*
|
||||
* @param size 申请空间大小
|
||||
* @param allocateHint 空间申请提示信息
|
||||
* @param extents 空间分配结果
|
||||
* @return uint64_t 已分配空间大小
|
||||
*/
|
||||
virtual uint64_t Allocate(uint32_t size, const AllocateHint& allocateHint,
|
||||
Extents* extents) = 0;
|
||||
|
||||
/**
|
||||
* @brief 释放空间
|
||||
*/
|
||||
virtual void Deallocate(const Extents& extents) = 0;
|
||||
|
||||
/**
|
||||
* @brief 标记对应空间已使用,初始化时使用
|
||||
*/
|
||||
virtual bool MarkUsed(const Extents& extents) = 0;
|
||||
|
||||
/**
|
||||
* @brief 标记对应空间可以使用,初始化时使用
|
||||
*/
|
||||
virtual bool MarkUsable(const Extents& extents) = 0;
|
||||
|
||||
/**
|
||||
* @brief 当前剩余空间
|
||||
*/
|
||||
virtual uint64_t TotalFree() const = 0;
|
||||
};
|
||||
```
|
||||
|
||||
MarkUsed 和 MarkFree 是持久化层调用,对分配器进行初始化。
|
||||
|
||||
## 后续版本需要解决的问题
|
||||
|
||||
- [ ] 空间分配信息如何重建
|
||||
- [ ] 高可用
|
||||
- [ ] 幂等请求处理
|
||||
- [ ] PB 级别情况下的性能
|
||||
|
|
@ -0,0 +1,178 @@
|
|||
# CurveFS 集成测试方案
|
||||
|
||||
## 目的
|
||||
|
||||
集成测试是在单元测试的基础上,将所有的软件单元按照概要设计规格说明的要求组装成模块、子系统进行测试。发生在单元测试之后和系统测试之前,用于验证不同模块间的接口调用、交互逻辑是否符合预期。
|
||||
|
||||
### 与单元测试区别
|
||||
|
||||
单元测试的关注点是一个范围很小的单元,通常是一个函数或一处关键逻辑。对于完成这个函数功能所依赖的外部组件,比如文件系统、数据库、网络请求等,进行 mock 或用 fake 对象代替。而集成测试通过组合相互依赖的单元/模块得到子系统,针对子系统进行测试。
|
||||
|
||||
### 与集成测试区别
|
||||
|
||||
系统测试属于黑盒测试,重点关注系统整体功能是否正常,异常场景下系统的表现是否符合预期等。而集成测试属于白盒或灰盒测试,可以更细粒度的关注处理流程。同时,在不同的测试方法中,比如自底向上集成,可以只针对某一个具体的功能进行测试,而不需要搭建整个集群。
|
||||
|
||||
## 测试方法
|
||||
|
||||
集成测试主要有两种执行方式,一种是一次性将所有单元组装起来进行测试的“大爆炸”模式。另外一种是层层递进的模式,每次集成一个新的模块进行测试,直至所有模块都组装完成,这里的递进也分为两种:自底向上和自顶向下。
|
||||
|
||||
### 大爆炸集成方法
|
||||
|
||||
在完成所有模块开发及单元测试后,将所有单元/模块组装到一起进行测试。这种方法与系统测试是有区别的,这里集成之后的测试重点还是各个子系统的功能,以及不同子系统之间的交互。而系统测试是整体功能的测试。
|
||||
|
||||

|
||||
|
||||
### 自顶向下集成方法
|
||||
|
||||
从上到下逐步测试单元/模块的集成。首先测试较高级别的模块,然后再测试和集成较低级别的模块,以检查软件功能。对于测试中未完成的底层模块,通过编写测试桩(stub)替代。
|
||||
|
||||

|
||||
|
||||
### 自底向上集成方法
|
||||
|
||||
单元/模块从底层到顶层,一步一步地进行测试,直到所有级别的单元/模块都集成在一起并作为一个整体进行测试。同样,在部分上层模块未完成开发是,需要编写测试驱动(driver)来启动测试。
|
||||
|
||||

|
||||
|
||||
## 测试内容
|
||||
|
||||
### 功能性测试
|
||||
|
||||
根据详细设计中模块提供的功能进行完备的测试。在做功能测试前需要设计充分的测试用例,考虑各种系统状态和参数输入下模块依然能够正常工作,并返回预期的结果。
|
||||
|
||||
在这一过程中,需要有一种科学的用例设计方法,依据这种方法可以充分考虑各种输入场景,保证测试功能不被遗漏,同时能够以尽可能少的用例和执行步骤完成测试过程。
|
||||
|
||||
### 异常测试
|
||||
|
||||
异常测试是有别于功能测试和性能测试又一种测试类型,通过异常测试,可以发现由于系统异常、依赖服务异常、应用本身异常等原因引起的性能瓶颈,提高系统的稳定性。
|
||||
|
||||
常见的异常如:磁盘错误、网络错误、数据出错、程序重启等等。
|
||||
|
||||
### 压力测试
|
||||
|
||||
功能性的测试更多是单线程的测试,还需要在并发场景下对模块进行测试,观察在并发场景下是否能够正常工作,逻辑或者数据是否会出现错误。
|
||||
|
||||
考虑并发度时,使用较大的压力来测试,例如平时使用的2倍、10倍或更大的压力。
|
||||
|
||||
## CurveBS 集成测试方案
|
||||
|
||||
前面说到,集成测试通常是在单元测试之后,系统测试之前完成,用于发现单元/模块组合过程中的问题。当前 FS 已经发布了两个 beta 版本,都经过 QA 完整的测试。结合之前 CurveBS 集成测试的经验,完成所有模块的集成测试,包括用例设计、代码编写、review 等,可能需要花费数周的时间。所以,整体来看,补充现有模块的集成测试收益不是很明显。但是,后续新特性的开发、bug 修复等,可以添加相应的集成测试。
|
||||
|
||||
### 新增功能集成测试
|
||||
|
||||
在当前的开发流程中,开发人员完成方案设计、代码开发、单元测试及代码 review、通过持续集成测试后,就可以将代码合入仓库。
|
||||
|
||||
在版本提测前,开发人员编写对应功能的测试用例,然后进行集群部署、人工测试、bug 修复,待所有功能完成自测后,提交版本给 QA 进行系统测试。单个功能的测试工作量在两个工作日左右。**这里的功能自测,其实就是对应功能的集成测试**。如果能在开发过程中,以集成测试代码的形式完成上述过程,可能会加快版本发布的流程。
|
||||
|
||||
### Bug 修复集成测试
|
||||
|
||||
以 [kudu](https://github.com/apache/kudu) 为例,[集成测试](https://github.com/apache/kudu/tree/master/src/kudu/integration-tests)中除了基本功能和压力测试外,包含了大量的针对 bug 的回归测试(100多个)。例如
|
||||
|
||||
```cpp
|
||||
// Regression test for KUDU-1551: if the tserver crashes after preallocating a segment
|
||||
// but before writing its header, the TS would previously crash on restart.
|
||||
// Instead, it should ignore the uninitialized segment.
|
||||
TEST_P(TsRecoveryITest, TestCrashBeforeWriteLogSegmentHeader) {
|
||||
NO_FATALS(StartClusterOneTs({
|
||||
"--log_segment_size_mb=1",
|
||||
"--log_compression_codec=NO_COMPRESSION"
|
||||
}));
|
||||
TestWorkload work(cluster_.get());
|
||||
work.set_num_replicas(1);
|
||||
work.set_write_timeout_millis(1000);
|
||||
work.set_timeout_allowed(true);
|
||||
work.set_payload_bytes(10000); // make logs roll without needing lots of ops.
|
||||
work.Setup();
|
||||
|
||||
// Enable the fault point after creating the table, but before writing any data.
|
||||
// Otherwise, we'd crash during creation of the tablet.
|
||||
ASSERT_OK(cluster_->SetFlag(cluster_->tablet_server(0),
|
||||
"fault_crash_before_write_log_segment_header", "0.9"));
|
||||
work.Start();
|
||||
|
||||
// Wait for the process to crash during log roll.
|
||||
ASSERT_OK(cluster_->tablet_server(0)->WaitForInjectedCrash(MonoDelta::FromSeconds(60)));
|
||||
work.StopAndJoin();
|
||||
|
||||
cluster_->tablet_server(0)->Shutdown();
|
||||
ignore_result(cluster_->tablet_server(0)->Restart());
|
||||
|
||||
ClusterVerifier v(cluster_.get());
|
||||
NO_FATALS(v.CheckRowCount(work.table_name(),
|
||||
ClusterVerifier::AT_LEAST,
|
||||
work.rows_inserted()));
|
||||
}
|
||||
```
|
||||
|
||||
用例描述了该测试用相应的 [issue](https://issues.apache.org/jira/browse/KUDU-1551)、触发场景等,集成测试用例与修复代码一起提交合入,以保证功能修复的正确性,也可以确保后续的代码修改不会对这次的修复产生影响。
|
||||
|
||||
同时,当前 FS 也加入了[自动化测试](https://github.com/opencurve/curve/blob/fs/robot/curve_fs_robot.txt),覆盖了基本功能、常见异常、数据一致性等测试。所以,在添加集成测试用例时,也需要考虑自动化测试是否可以覆盖需要测试的场景,如果可以,就不需要添加相应的集成测试用例。反之,如果没有覆盖,或者覆盖不完全,则需要添加。
|
||||
|
||||
所以,综合以上的讨论,是否添加相应的集成测试,可以参考如下的判断:
|
||||
|
||||
- 是否需要多模块协作
|
||||
- 如果是,则需要再完成单元测试后,添加整体功能的集成测试。
|
||||
- 已有的自动化测试是否能够覆盖主要流程
|
||||
- 如果不能,则可以根据工作量、改动情况再决定添加自动化测试或集成测试。
|
||||
- 是否影响现有的功能
|
||||
- 如果是,则同时需要对现有功能添加相应的集成测试,保证现有功能的正确性。
|
||||
- ...
|
||||
|
||||
## 集成测试用例设计方法
|
||||
|
||||
测试用例是为验证程序是否符合特定系统需求而开发的测试输入、执行条件和预期结果的集合,其组织性、功能覆盖性、重复性的特点能够保证测试功能不被遗漏。由于测试用例往往涉及多重选择、循环嵌套,不同的路径数目可能非常大,所以必须精心设计使之能达到最佳的测试效果。
|
||||
|
||||
### 设计原则
|
||||
|
||||
这里首先需要考虑从什么样的角度去思考用例的设计,一种是以接口的角度,根据接口划分来设计用例;另一种是以使用场景的角度来设计。
|
||||
|
||||
#### 根据接口来设计用例
|
||||
|
||||
优点:可以根据输入不同的接口参数,产生不同的预期来设计用例,考虑比较完整地考虑各种调用情况。
|
||||
|
||||
缺点:根据接口提供的功能来设计用例,有时候难以考虑一些特殊场景下功能上的缺陷。
|
||||
|
||||
#### 根据场景来设计用例
|
||||
|
||||
优点:用例审核者能比较容易理解各组用例出现的场景;可以测试出功能设计之初未考虑到的一些场景。
|
||||
|
||||
缺点:无法很好地证明用例覆盖是否充分。
|
||||
|
||||
两种方式各有优缺点,可以考虑将两种方式结合起来,用接口方式来设计用例,然后用场景方式来组织用例的执行序列。
|
||||
|
||||
根据接口进行设计可以比较清晰地评估用例覆盖是否完全,然后以场景方式来执行这些用例,对执行过的用例就打钩记录;这样可以很清楚的知道哪些用例执行了哪些没执行,帮助发现没有想到的场景;
|
||||
|
||||
此外还能通过特殊的场景帮助发现接口功能是否考虑充分,两种方式可以互补帮助发现更多问题。
|
||||
|
||||
### 用例模板
|
||||
|
||||
用例设计可以遵循GWT(Given-When-Then)的模式来写,Given表示给定的前提条件,When表示要发生的操作,Then表示预期的结果。
|
||||
|
||||
如果要覆盖所有的用例,那么势必要列出所有的前提条件,然后在特定的前提下,需要列出所有可能发生的操作。
|
||||
|
||||
| 编号 | Given | When | Than | 备注 | 是否执行 |
|
||||
| ---- | ---- | ---- | ----- | --- | ------ |
|
||||
| 1 | x | x | x | x | x |
|
||||
|
||||
用例设计可以参考之前 CurveBS Datastore 的集成测试用例(https://github.com/opencurve/curve/blob/master/docs/cn/quality-integration-example.md)。
|
||||
|
||||
## 模糊测试
|
||||
|
||||
无论是单元测试、集成测试还是自动化测试,都需要开发或测试人员设计测试用例并编写相应的代码进行测试,同时,这些测试通常只对主要流程进行针对性的测试。而模糊测试(Fuzzing、Fuzz Testing)是一种自动化的测试方法,可以向系统提供非法、超出预期或随机的数据,并监控系统行为是否异常,以发现可能的错误。
|
||||
|
||||
[OSS-Fuzz](https://github.com/google/oss-fuzz) 是当前比较知名的模糊测试框架之一,由 Google 开源,支持多种开发语言,截止到 2022 年 1 月份,已发现 550 个开源软件中的 36,000 多个 bug。
|
||||
|
||||
TODO:
|
||||
|
||||
- [ ] 如何使用?
|
||||
- [ ] 是否适用?
|
||||
|
||||
## 参考
|
||||
|
||||
1 [The Differences Between Unit Testing, Integration Testing And Functional Testing](https://www.softwaretestinghelp.com/the-difference-between-unit-integration-and-functional-testing/)
|
||||
|
||||
2 [软件测试入门系列之十:集成测试](https://zhuanlan.zhihu.com/p/354967307)
|
||||
|
||||
3 [Fuzzing](https://en.wikipedia.org/wiki/Fuzzing)
|
||||
|
||||
4 [Announcing OSS-Fuzz: Continuous fuzzing for open source software](https://opensource.googleblog.com/2016/12/announcing-oss-fuzz-continuous-fuzzing.html)
|
||||
|
|
@ -0,0 +1,319 @@
|
|||
# 背景
|
||||
|
||||
如果Curvefs client在写底层存储的时候是直接写入远端对象存储,那么由于写远端时延相对会较高,所以为了提升性能,引入了写本地缓存盘方案。也即要写底层存储时,先把数据写到本地缓存硬盘,然后再把本地缓存硬盘中的数据异步上传到远端对象存储。
|
||||
|
||||
但是被当作缓存盘的盘有可能是机器的系统盘(如果系统盘io打满了,会影响整个服务器的服务)或者其他一些原因需要对缓存盘的写入或者读取进行限速,因此需要对缓存盘进行限速。
|
||||
|
||||
同时,希望把该限速方案设计为一个可供curve其他模块使用的通用模块。
|
||||
|
||||
# 业界方案调研
|
||||
|
||||
一般来说,限速的话有throttle和qos两种,qos的话就是严格控制每秒钟的iops和带宽,throttle的话就是根据后端处理能力调节前端的调用,每秒钟的iops和带宽是不固定的。显然,对于curvefs本地盘限速来说,qos的设计是更符合要求的。
|
||||
|
||||
在这里先简单介绍下throttle,throttle有两种:一种是严格的throttle,也就是达到限流水位后,只有释放一个后端的能力,才能有一个前端的调用。另外一种是自适应throttle,,也即根据当前的水位来调节,有低水位阈值和高水位阈值。如果当前值处于低水位阈值之下,那么无需调节,如果当前的水位处于低水位和高水位之间,轻微调节,让进程等待少量的时间后再处理请求,如果当前的水位比高水位阈值还高,那么就让进程程多等待一些时间。Ceph中的BackoffThrottle算法就是这样一种自适应throttle算法。其限流思路如下图:
|
||||

|
||||
|
||||
y轴是延时时间,x轴代表当前水位。
|
||||
|
||||
low代表低水位阈值,high代表高水位阈值
|
||||
|
||||
QoS是Quality of Service的缩写,它起源于网络技术,用以解决网络延迟和阻塞等问题,能够为指定的网络通信提供更好的服务能力。有两种主流的QOS算法,漏桶算法和令牌桶算法。
|
||||
|
||||
漏桶算法就是将请求放入桶中,然后始终以一个固定的速率从桶中取出请求来处理,当桶中等待的请求数超过上限后(桶的容量固定),后续的请求就不再加入桶中,而是执行拒绝策略(比如降级)。漏桶算法如下图。漏桶算法适用于需要以固定速率的场景,而在多数业务场景中,我们并不需要严格的速率,并且需要有一定的应对突发流量的能力,所以会使用令牌桶算法限流。漏桶算法如下图:
|
||||

|
||||
|
||||
令牌桶算法就是以固定速率生成令牌放入桶中,每个请求都需要从桶中获取令牌,没有获取到令牌的请求会被阻塞限流(桶中的令牌不够的时候),当令牌消耗速度小于生成的速度时,令牌桶内就会预存这些未消耗的令牌(直到桶的上限),当有突发流量进来时,可以直接从桶中取出令牌,而不会被限流,令牌桶算法如下图:
|
||||

|
||||
|
||||
对于curvefs缓存盘的限速,我们这里采用令牌桶算法。
|
||||
|
||||
# 设计方案
|
||||
|
||||
令牌桶的设计以及工作过程如下:
|
||||
|
||||
- 令牌根据时间匀速的产生令牌数量(这里假设是putTokens_),存入到令牌桶中,同时保证令牌桶中的令牌最大数量为maxTokens_。并且要保证putTokens_<=maxTokens_。另外令牌桶在初始化的时候,会分配一定数量的令牌数(putTokens_)。
|
||||
|
||||
|
||||
- 当线程要获取令牌消费时,这里假设消费者需要m个令牌,如果能获取得到足够的令牌,那么消费者获取成功;如果获取不到足够的令牌,那么就触发保护策略,消费者线程等待。
|
||||
|
||||
- 定时匀速向令牌桶丢putTokens_个令牌,如果有消费者在等待令牌,那么则唤醒消费者(当然,消费者要等到令牌桶中有足够的令牌个数了才会被唤醒)
|
||||
|
||||
关键数据结构如下:
|
||||
```
|
||||
// 令牌桶管理模块
|
||||
class TokenBucketControl {
|
||||
public:
|
||||
TokenBucketControl(const std::string& n, uint64_t maxToken, uint64_t putTokens) : name_(n),
|
||||
putTokens_(putTokens), tokenthrottle_(maxToken) {
|
||||
if (maxToken < putTokens) {
|
||||
tokenthrottle_.SetMaxTokens(putTokens);
|
||||
}
|
||||
}
|
||||
~TokenBucketControl() {}
|
||||
// 添加令牌
|
||||
void AddTokens();
|
||||
// 消费者线程获取令牌
|
||||
void GetTokens(uint64_t reqToken);
|
||||
private:
|
||||
// 调用模块的名称(该限流模块是通用模块,可以被curve中不同模块服务)
|
||||
const std::string name_;
|
||||
|
||||
// 向令牌桶中丢令牌的线程
|
||||
curve::common::Thread backEndThread_;
|
||||
curve::common::Atomic<bool> isRunning_;
|
||||
curve::common::InterruptibleSleeper sleeper_;
|
||||
uint32_t PutTokensIntervalMs_;
|
||||
// 单位时间
|
||||
uint32_t splitWindowsSec_;
|
||||
// 单位之间被分拆成splitWindowsNums_个窗口单元
|
||||
uint32_t splitWindowsNums_;
|
||||
int PutTokensRun();
|
||||
int PutTokensStop();
|
||||
virtual int PutTokensFun();
|
||||
|
||||
// 令牌桶
|
||||
class TokenBucket {
|
||||
private:
|
||||
// 令牌桶中可以存放的最大令牌数
|
||||
uint64_t maxTokens_;
|
||||
// 令牌桶中当前令牌个数
|
||||
uint64_t remainTokens_;
|
||||
public:
|
||||
TokenBucket(uint64_t m)
|
||||
: remainTokens_(m), maxTokens_(m) {
|
||||
}
|
||||
~TokenBucket() {}
|
||||
// 获取令牌
|
||||
bool Get(uint64_t tokens);
|
||||
// 放置令牌
|
||||
void Put(uint64_t tokens);
|
||||
void SetMaxTokens(uint64_t tokens) { maxTokens_ = tokens;}
|
||||
};
|
||||
TokenBucket tokenthrottle_;
|
||||
|
||||
// 每次放置令牌的数量
|
||||
uint64_t putTokens_;
|
||||
curve::common::Mutex mtx_;
|
||||
|
||||
// 消费者线程等待 && 唤醒消费者
|
||||
std::mutex limitMtx_;
|
||||
std::condition_variable limitCond_;
|
||||
void WaitTokenLimit() {
|
||||
std::unique_lock<std::mutex> lk(limitMtx_);
|
||||
limitCond_.wait(lk);
|
||||
}
|
||||
void SignalTokenLimit() {
|
||||
std::lock_guard<std::mutex> lk(limitMtx_);
|
||||
limitCond_.notify_all();
|
||||
}
|
||||
|
||||
};
|
||||
```
|
||||
## 要点说明
|
||||
|
||||
- 流量不均匀(突然爆发)问题
|
||||
|
||||
比如我们每秒限定的iops是1000,那么有可能在前零点几秒一下就把这1000个iops对应的令牌全部用完了,造成流量突然爆发,但是该一秒钟内其余时间段就获取不到令牌了。
|
||||
|
||||
那么解决这个问题的办法就是把这一秒的时间分成splitWindowsNums_个时间窗口,每个时间窗口投放1000/splitWindowsNums_个令牌。
|
||||
|
||||
- 峰值流量不能持续(1秒)问题
|
||||
|
||||
比如我们设置的令牌的每秒放置量是putTokens_,然后设置的最大令牌数量是maxTokens_,那么这maxTokens_有可能在1秒钟就全部被拿完了。
|
||||
|
||||
那么想解决这个问题(比如我们期望峰值流量持续burst_length 秒)可以根据如下的思路:
|
||||
|
||||
1. 使用两个桶
|
||||
2. 大桶的令牌个数最多为`maxTokens_*burst_length`,小桶的令牌个数最多为maxTokens_。
|
||||
3. 放置令牌:每秒向大桶放putTokens_个令牌;向小桶放`burst_ratio*putTokens`个令牌,也即maxTokens_个令牌(`burst_ratio=maxTokens_/putTokens`,burst_ratio表示峰值是平均值的几倍)。另外,小桶的令牌个数不大于maxTokens_,同时要保证小桶的令牌个数不超过大桶中的令牌个数(这一点很关键,因为大桶中令牌个数最多为`maxTokens_*burst_length`,当峰值流量来时,由于小桶每秒都能获得峰值个令牌,那么可以认为小桶是不受限的,并且由于小桶中的令牌个数不能超过大桶中的令牌个数,这样就保证了峰值流量的持续时间是burst_length----其实这里还会有误差,因为取令牌的同时,还在向里面放令牌,所以峰值时间比实际的时间要长,如果峰值跟均值差距不大,那么这个误差值可能很大,这个解决办法见下面的措施)。
|
||||
4. 取令牌:若峰值来了,要取maxTokens_个令牌,那么每次分别从小桶和大桶中取maxTokens_个令牌。那么也就是每秒可取maxTokens_个令牌(也即峰值)。
|
||||
|
||||
- 峰值流量时间超出预期的问题
|
||||
|
||||
举个例子,比如putTokens_是80,maxTokens是100,我们期望的峰值时间burst_length是60s。那么根据上述所说大桶可放置的最大令牌个数便是maxTokens*burst_length。当由于峰值来临的时候,每秒拿的令牌个数是maxTokens,但与此同时每秒钟又向大桶放入了putTokens_个令牌。所以如果大桶可放的最大令牌个数是`maxTokens*burst_lengt`h。那么他实际可保持的峰值时间是:`(maxTokens * burst_length)/(maxTokens-putTokens_)=(100 * 60)/(100-80) = 300`。故而实际的峰值时间远超期待的峰值时间。
|
||||
|
||||
那么要想控制峰值时间的准确性,可以通过修改大桶可放的最大令牌个数来优化。假设大桶可放置的最大令牌个数为max。那么可以通过如下公式计算得到: `max/(maxTokens-putTokens_)=burst_length`, 所以`max=burst_length*(maxTokens-putTokens_)`。那么再进行一些小调整(来源于https://github.com/ceph/ceph/pull/35138):`max = maxTokens + (maxTokens - putTokens_)* (burst_length - 1)`
|
||||
|
||||
- 瞬时流量可能超过期望峰值的问题
|
||||
|
||||
如果一段时间没有io请求过来,那么由于每秒都会向小令牌桶放maxTokens个令牌,但是又没有进程取令牌,所以此时令牌桶便有可能被装满了。因此当流量来时,第一秒中的有效令牌个数便是令牌桶里当前所有的令牌再加上这一秒放进去的令牌,也即最大可能时2*maxTokens,也即导致瞬时流量极大。
|
||||
|
||||
假设maxTokens是500。那么如果一段时间没有io请求过来,当突然有io过来时,此时小令牌桶可能已经满了,有maxTokens个令牌,那么第一秒的瞬时流量有可能达到500+500。
|
||||
|
||||
这个问题的解决方案同下面漏桶一小节描述。
|
||||
|
||||
## 使用
|
||||
|
||||
这里以磁盘缓存管理类`DiskCacheManager`中的使用为例
|
||||
|
||||
- 初始化
|
||||
```
|
||||
class DiskCacheManager {
|
||||
// 分别定义写带宽和iops限流管理类
|
||||
std::shared_ptr<TokenBucketControl> writeBwControl_;
|
||||
std::shared_ptr<TokenBucketControl> writeIopsControl_;
|
||||
}
|
||||
|
||||
DiskCacheManager::DiskCacheManager(...) {
|
||||
// 初始化,设定相应限流参数
|
||||
writeBwControl_ =
|
||||
std::make_shared<TokenBucketControl>("diskcache_wb", diskCache.avgPutFlushBytes, diskCache.maxFlushBytes);
|
||||
writeIopsControl_ =
|
||||
std::make_shared<TokenBucketControl>("diskcache_wi", diskCache.avgPutFlushIops, diskCache.maxFlushIops);
|
||||
// 启动限流
|
||||
writeBwControl_->PutTokensRun();
|
||||
writeIopsControl_->PutTokensRun();
|
||||
}
|
||||
|
||||
```
|
||||
|
||||
- 使用
|
||||
```
|
||||
void DiskCacheManager::GetTokens(bool isWrite, uint64_t len) {
|
||||
if (isWrite) {
|
||||
writeBwControl_->GetTokens(len);
|
||||
writeIopsControl_->GetTokens(1);
|
||||
} else {
|
||||
readBwControl_->GetTokens(len);
|
||||
readIopsControl_->GetTokens(1);
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
int DiskCacheManager::WriteDiskFile(const std::string fileName, const char *buf,uint64_t length, bool force) {
|
||||
// 获取限流指标,若未达到限流值,则继续运行;反之,则wait等待唤醒
|
||||
GetTokens(length);
|
||||
|
||||
// do other things
|
||||
}
|
||||
```
|
||||
|
||||
# poc测试结果
|
||||
poc测试中对缓存盘进行限流,限流带宽为100MB(104857600)
|
||||
|
||||

|
||||
|
||||
# 漏桶
|
||||
## 设计方案
|
||||
|
||||
漏桶的设计以及工作过程如下:
|
||||
|
||||
- 漏桶匀速漏水(这里假设漏水速率是avg),也即每秒可以处理数量为avg的任务。另外桶的容量也为avg
|
||||
|
||||
|
||||
- 当任务来临时,往桶内加对应的水,桶内水的总量不能超过桶的容量。桶内还可以添加的容量就是桶的容量减去桶的当前水位。
|
||||
|
||||
- 如果桶内可用容量小于本次任务请求的数量,那么便等待漏水。如果请求的数量为N*avg(且此时桶已满),那么就需要漏N次水才能唤醒本次任务
|
||||
|
||||
## 要点说明
|
||||
|
||||
这里的几个要点其实跟令牌桶算法那里描述的是差不多的
|
||||
|
||||
- 流量不均匀(突然爆发)问题
|
||||
|
||||
也是把每秒拆分成N个窗口,每个窗口的漏水量为avg/N
|
||||
|
||||
- 峰值流量不能持续(1秒)问题
|
||||
|
||||
其解决思路与上述令牌桶中的思路也类似。想解决这个问题(比如我们期望峰值流量持续burst_length 秒)可以根据如下的思路:
|
||||
|
||||
1. 使用两个桶
|
||||
2. 大桶的容量为burst(峰值)*burst_length,小桶的容量为burst。
|
||||
3. 如果漏桶内有水,那么漏桶匀速漏水。大桶每秒漏水量为avg;小桶每秒漏水量为burst,由于小桶每秒漏水量就是峰值burst,所以小桶每秒后都是空的(所以可用容量一直是burst)。我们可以通过取有效容量为大桶和小桶中的可用容量小者,由于大桶的容量是burst(峰值)*burst_length,所以大桶在burst_length秒的时间会被装满(当然,与令牌桶中描述的类似,这里还会有误差,因为添水的同时,大桶也一直在漏水)
|
||||
|
||||
- 峰值流量时间超出预期的问题
|
||||
|
||||
同令牌同一小节描述类似,也是通过控制大桶的容量来解决这个问题
|
||||
|
||||
- 瞬时流量可能超过期望峰值的问题
|
||||
|
||||
如果一段时间没有io请求过来,那么由于小桶每秒都在漏水,但是又没有进程添水,所以此时小桶便有可能已经彻底为空了,其有效容量为burst。因此当流量来时,第一秒中的有效容量便是小桶当前的可用容量加第一秒漏水两,也即最大可能是2*burst,也即导致瞬时流量极大。假设burst是500,那么瞬时流量有可能达到500+500。
|
||||
|
||||
解决思路可以是当小桶中有连续两个时间周期都没有水,那么便初始化小桶的容量为一个时间周期的漏水量, 大体如下:
|
||||
|
||||
```
|
||||
class LeakyBucket {
|
||||
struct Bucket {
|
||||
bool initial = true;
|
||||
uint32_t noBurstLevelTimes = 0;
|
||||
bthread::Mutex mtx_;
|
||||
}
|
||||
}
|
||||
|
||||
void LeakyBucket::Bucket::Leak(uint64_t intervalUs) {
|
||||
double leak = static_cast<double>(avg) * intervalUs /
|
||||
TimeUtility::MicroSecondsPerSecond;
|
||||
level = std::max(level - leak, 0.0);
|
||||
if (burst > 0) {
|
||||
if (burstLevel == 0) {
|
||||
++noBurstLevelTimes;
|
||||
if (noBurstLevelTimes > 2) {
|
||||
std::lock_guard<bthread::Mutex> lk(mtx_);
|
||||
// init
|
||||
burstLevel = burst - (static_cast<double>(burst) * FLAGS_bucketLeakIntervalMs /
|
||||
TimeUtility::MicroSecondsPerSecond);;
|
||||
noBurstLevelTimes = 0;
|
||||
initial = true;
|
||||
}
|
||||
} else {
|
||||
noBurstLevelTimes = 0;
|
||||
}
|
||||
// if initial is true, return direct
|
||||
if (initial) {
|
||||
return;
|
||||
}
|
||||
leak = static_cast<double>(burst) * intervalUs /
|
||||
TimeUtility::MicroSecondsPerSecond;
|
||||
burstLevel = std::max(burstLevel - leak, 0.0);
|
||||
}
|
||||
}
|
||||
void LeakyBucket::Bucket::Reset(uint64_t avg, uint64_t burst,
|
||||
uint64_t burstLength) {
|
||||
// init burstLevel
|
||||
this->burstLevel = burst - (static_cast<double>(burst) * FLAGS_bucketLeakIntervalMs /
|
||||
TimeUtility::MicroSecondsPerSecond);
|
||||
}
|
||||
```
|
||||
|
||||
## 使用
|
||||
|
||||
这里依然以这里以磁盘缓存管理类DiskCacheManager中的使用为例:
|
||||
|
||||
- 初始化
|
||||
```
|
||||
class DiskCacheManager {
|
||||
Throttle diskCacheThrottle_;
|
||||
}
|
||||
int DiskCacheManager::Init(...) {
|
||||
ReadWriteThrottleParams params;
|
||||
params.iopsWrite = ThrottleParams(option.diskCacheOpt.maxFlushNums, 0, 0);
|
||||
params.bpsWrite = ThrottleParams(option.diskCacheOpt.maxFlushBytes, 0, 0);
|
||||
params.iopsRead = ThrottleParams(option.diskCacheOpt.maxReadFileNums, 0, 0);
|
||||
params.bpsRead = ThrottleParams(option.diskCacheOpt.maxReadFileBytes, 0, 0);
|
||||
|
||||
diskCacheThrottle_.UpdateThrottleParams(params);
|
||||
}
|
||||
```
|
||||
- 使用
|
||||
|
||||
```
|
||||
int DiskCacheManager::WriteDiskFile(const std::string fileName, const char *buf,
|
||||
uint64_t length, bool force) {
|
||||
// write bps throttle
|
||||
diskCacheThrottle_.Add(false, length);
|
||||
// write iops throttle
|
||||
diskCacheThrottle_.Add(false, 1);
|
||||
int ret = cacheWrite_->WriteDiskFile(fileName, buf, length, force);
|
||||
if (ret > 0)
|
||||
AddDiskUsedBytes(ret);
|
||||
return ret;
|
||||
}
|
||||
```
|
||||
|
||||
# 参考
|
||||
[漏桶算法和令牌桶算法](https://juejin.cn/post/6961815018488725541)
|
||||
|
||||
[ceph qos设计](https://github.com/ceph/ceph/pull/17032)
|
||||
|
||||
|
||||
|
||||
|
After Width: | Height: | Size: 62 KiB |
|
After Width: | Height: | Size: 22 KiB |
|
After Width: | Height: | Size: 31 KiB |
|
After Width: | Height: | Size: 19 KiB |
|
After Width: | Height: | Size: 15 KiB |
|
After Width: | Height: | Size: 3.0 KiB |
|
After Width: | Height: | Size: 13 KiB |
|
After Width: | Height: | Size: 12 KiB |
|
After Width: | Height: | Size: 7.9 KiB |
|
After Width: | Height: | Size: 346 KiB |
|
After Width: | Height: | Size: 655 KiB |
|
After Width: | Height: | Size: 360 KiB |