Commit 593082c4 authored by zhujiashun's avatar zhujiashun

revived_from_all_failed: refine comments and UT

parent 57b60aa4
...@@ -129,7 +129,10 @@ struct ChannelOptions { ...@@ -129,7 +129,10 @@ struct ChannelOptions {
// Default: "" // Default: ""
std::string connection_group; std::string connection_group;
// TODO(zhujiashun) // Customize the revive policy after all servers are shutdown. The
// interface is defined in src/brpc/revive_policy.h
// This object is NOT owned by channel and should remain valid when
// channel is used.
// Default: NULL // Default: NULL
RevivePolicy* revive_policy; RevivePolicy* revive_policy;
......
...@@ -302,6 +302,20 @@ int ConsistentHashingLoadBalancer::SelectServer( ...@@ -302,6 +302,20 @@ int ConsistentHashingLoadBalancer::SelectServer(
if (s->empty()) { if (s->empty()) {
return ENODATA; return ENODATA;
} }
RevivePolicy* rp = in.revive_policy;
if (rp) {
std::set<ServerId> server_list;
for (auto server: *s) {
server_list.insert(server.server_sock);
}
std::vector<ServerId> server_list_distinct(
server_list.begin(), server_list.end());
if (rp->DoReject(server_list_distinct)) {
return EREJECT;
}
rp->StopRevivingIfNecessary();
}
std::vector<Node>::const_iterator choice = std::vector<Node>::const_iterator choice =
std::lower_bound(s->begin(), s->end(), (uint32_t)in.request_code); std::lower_bound(s->begin(), s->end(), (uint32_t)in.request_code);
if (choice == s->end()) { if (choice == s->end()) {
...@@ -319,6 +333,9 @@ int ConsistentHashingLoadBalancer::SelectServer( ...@@ -319,6 +333,9 @@ int ConsistentHashingLoadBalancer::SelectServer(
} }
} }
} }
if (rp) {
rp->StartReviving();
}
return EHOSTDOWN; return EHOSTDOWN;
} }
......
...@@ -22,7 +22,7 @@ ...@@ -22,7 +22,7 @@
#include "brpc/socket.h" #include "brpc/socket.h"
#include "brpc/reloadable_flags.h" #include "brpc/reloadable_flags.h"
#include "brpc/policy/locality_aware_load_balancer.h" #include "brpc/policy/locality_aware_load_balancer.h"
#include "brpc/revive_policy.h"
namespace brpc { namespace brpc {
namespace policy { namespace policy {
...@@ -270,7 +270,6 @@ int LocalityAwareLoadBalancer::SelectServer(const SelectIn& in, SelectOut* out) ...@@ -270,7 +270,6 @@ int LocalityAwareLoadBalancer::SelectServer(const SelectIn& in, SelectOut* out)
if (n == 0) { if (n == 0) {
return ENODATA; return ENODATA;
} }
size_t ntry = 0; size_t ntry = 0;
size_t nloop = 0; size_t nloop = 0;
int64_t total = _total.load(butil::memory_order_relaxed); int64_t total = _total.load(butil::memory_order_relaxed);
......
...@@ -110,9 +110,12 @@ int RandomizedLoadBalancer::SelectServer(const SelectIn& in, SelectOut* out) { ...@@ -110,9 +110,12 @@ int RandomizedLoadBalancer::SelectServer(const SelectIn& in, SelectOut* out) {
if (n == 0) { if (n == 0) {
return ENODATA; return ENODATA;
} }
if (in.revive_policy && RevivePolicy* rp = in.revive_policy;
(in.revive_policy)->RejectDuringReviving(s->server_list)) { if (rp) {
return EREJECT; if (rp->DoReject(s->server_list)) {
return EREJECT;
}
rp->StopRevivingIfNecessary();
} }
uint32_t stride = 0; uint32_t stride = 0;
size_t offset = butil::fast_rand_less_than(n); size_t offset = butil::fast_rand_less_than(n);
...@@ -132,8 +135,8 @@ int RandomizedLoadBalancer::SelectServer(const SelectIn& in, SelectOut* out) { ...@@ -132,8 +135,8 @@ int RandomizedLoadBalancer::SelectServer(const SelectIn& in, SelectOut* out) {
// this failed server won't be visited again inside for // this failed server won't be visited again inside for
offset = (offset + stride) % n; offset = (offset + stride) % n;
} }
if (in.revive_policy) { if (rp) {
in.revive_policy->StartRevive(); rp->StartReviving();
} }
// After we traversed the whole server list, there is still no // After we traversed the whole server list, there is still no
// available server // available server
......
...@@ -31,12 +31,6 @@ namespace policy { ...@@ -31,12 +31,6 @@ namespace policy {
// than RoundRobinLoadBalancer. // than RoundRobinLoadBalancer.
class RandomizedLoadBalancer : public LoadBalancer { class RandomizedLoadBalancer : public LoadBalancer {
public: public:
RandomizedLoadBalancer()
: _reviving(false)
//TODO(zhujiashun)
, _minimum_working_instances(2)
, _last_usable(0)
, _last_usable_change_time_ms(0) {}
bool AddServer(const ServerId& id); bool AddServer(const ServerId& id);
bool RemoveServer(const ServerId& id); bool RemoveServer(const ServerId& id);
size_t AddServersInBatch(const std::vector<ServerId>& servers); size_t AddServersInBatch(const std::vector<ServerId>& servers);
...@@ -57,13 +51,6 @@ private: ...@@ -57,13 +51,6 @@ private:
static size_t BatchRemove(Servers& bg, const std::vector<ServerId>& servers); static size_t BatchRemove(Servers& bg, const std::vector<ServerId>& servers);
butil::DoublyBufferedData<Servers> _db_servers; butil::DoublyBufferedData<Servers> _db_servers;
bool _reviving;
int64_t _minimum_working_instances;
// TODO(zhujiashun): remove mutex
butil::Mutex _mutex;
int64_t _last_usable;
int64_t _last_usable_change_time_ms;
}; };
} // namespace policy } // namespace policy
......
...@@ -110,9 +110,12 @@ int RoundRobinLoadBalancer::SelectServer(const SelectIn& in, SelectOut* out) { ...@@ -110,9 +110,12 @@ int RoundRobinLoadBalancer::SelectServer(const SelectIn& in, SelectOut* out) {
if (n == 0) { if (n == 0) {
return ENODATA; return ENODATA;
} }
if (in.revive_policy && RevivePolicy* rp = in.revive_policy;
(in.revive_policy)->RejectDuringReviving(s->server_list)) { if (rp) {
return EREJECT; if (rp->DoReject(s->server_list)) {
return EREJECT;
}
rp->StopRevivingIfNecessary();
} }
TLS tls = s.tls(); TLS tls = s.tls();
if (tls.stride == 0) { if (tls.stride == 0) {
...@@ -131,8 +134,8 @@ int RoundRobinLoadBalancer::SelectServer(const SelectIn& in, SelectOut* out) { ...@@ -131,8 +134,8 @@ int RoundRobinLoadBalancer::SelectServer(const SelectIn& in, SelectOut* out) {
return 0; return 0;
} }
} }
if (in.revive_policy) { if (rp) {
in.revive_policy->StartRevive(); rp->StartReviving();
} }
s.tls() = tls; s.tls() = tls;
return EHOSTDOWN; return EHOSTDOWN;
......
...@@ -20,6 +20,7 @@ ...@@ -20,6 +20,7 @@
#include "brpc/socket.h" #include "brpc/socket.h"
#include "brpc/policy/weighted_round_robin_load_balancer.h" #include "brpc/policy/weighted_round_robin_load_balancer.h"
#include "butil/strings/string_number_conversions.h" #include "butil/strings/string_number_conversions.h"
#include "brpc/revive_policy.h"
namespace { namespace {
...@@ -157,6 +158,18 @@ int WeightedRoundRobinLoadBalancer::SelectServer(const SelectIn& in, SelectOut* ...@@ -157,6 +158,18 @@ int WeightedRoundRobinLoadBalancer::SelectServer(const SelectIn& in, SelectOut*
if (s->server_list.empty()) { if (s->server_list.empty()) {
return ENODATA; return ENODATA;
} }
RevivePolicy* rp = in.revive_policy;
if (rp) {
std::vector<ServerId> server_list;
server_list.reserve(s->server_list.size());
for (auto server: s->server_list) {
server_list.emplace_back(server.id);
}
if (rp->DoReject(server_list)) {
return EREJECT;
}
rp->StopRevivingIfNecessary();
}
TLS& tls = s.tls(); TLS& tls = s.tls();
if (tls.IsNeededCaculateNewStride(s->weight_sum, s->server_list.size())) { if (tls.IsNeededCaculateNewStride(s->weight_sum, s->server_list.size())) {
if (tls.stride == 0) { if (tls.stride == 0) {
...@@ -198,6 +211,9 @@ int WeightedRoundRobinLoadBalancer::SelectServer(const SelectIn& in, SelectOut* ...@@ -198,6 +211,9 @@ int WeightedRoundRobinLoadBalancer::SelectServer(const SelectIn& in, SelectOut*
tls_temp.remain_server = tls.remain_server; tls_temp.remain_server = tls.remain_server;
} }
} }
if (rp) {
rp->StartReviving();
}
return EHOSTDOWN; return EHOSTDOWN;
} }
......
...@@ -34,45 +34,55 @@ DefaultRevivePolicy::DefaultRevivePolicy( ...@@ -34,45 +34,55 @@ DefaultRevivePolicy::DefaultRevivePolicy(
, _hold_time_ms(hold_time_ms) { } , _hold_time_ms(hold_time_ms) { }
void DefaultRevivePolicy::StartRevive() { void DefaultRevivePolicy::StartReviving() {
std::unique_lock<butil::Mutex> mu(_mutex);
_reviving = true; _reviving = true;
} }
bool DefaultRevivePolicy::RejectDuringReviving( void DefaultRevivePolicy::StopRevivingIfNecessary() {
const std::vector<ServerId>& server_list) { int64_t now_ms = butil::gettimeofday_ms();
if (!_reviving) { {
return false; std::unique_lock<butil::Mutex> mu(_mutex);
if (_last_usable_change_time_ms != 0 && _last_usable != 0 &&
(now_ms - _last_usable_change_time_ms > _hold_time_ms)) {
_reviving = false;
_last_usable_change_time_ms = 0;
}
}
return;
}
bool DefaultRevivePolicy::DoReject(const std::vector<ServerId>& server_list) {
{
std::unique_lock<butil::Mutex> mu(_mutex);
if (!_reviving) {
mu.unlock();
return false;
}
} }
size_t n = server_list.size(); size_t n = server_list.size();
int usable = 0; int usable = 0;
// TODO(zhujiashun): optimize looking process
SocketUniquePtr ptr; SocketUniquePtr ptr;
// TODO(zhujiashun): optimize O(N)
for (size_t i = 0; i < n; ++i) { for (size_t i = 0; i < n; ++i) {
if (Socket::Address(server_list[i].id, &ptr) == 0 if (Socket::Address(server_list[i].id, &ptr) == 0
&& !ptr->IsLogOff()) { && !ptr->IsLogOff()) {
usable++; usable++;
} }
} }
std::unique_lock<butil::Mutex> mu(_mutex); int64_t now_ms = butil::gettimeofday_ms();
if (_last_usable_change_time_ms != 0 && usable != 0 && {
(butil::gettimeofday_ms() - _last_usable_change_time_ms > _hold_time_ms) std::unique_lock<butil::Mutex> mu(_mutex);
&& _last_usable == usable) {
_reviving = false;
_last_usable_change_time_ms = 0;
mu.unlock();
} else {
if (_last_usable != usable) { if (_last_usable != usable) {
_last_usable = usable; _last_usable = usable;
_last_usable_change_time_ms = butil::gettimeofday_ms(); _last_usable_change_time_ms = now_ms;
}
mu.unlock();
int rand = butil::fast_rand_less_than(_minimum_working_instances);
if (rand >= usable) {
return true;
} }
} }
int rand = butil::fast_rand_less_than(_minimum_working_instances);
if (rand >= usable) {
return true;
}
return false; return false;
} }
} // namespace brpc } // namespace brpc
...@@ -23,20 +23,40 @@ ...@@ -23,20 +23,40 @@
namespace brpc { namespace brpc {
class ServerId; class ServerId;
// After all servers are shutdown and health check happens, servers are
// online one by one. Once one server is up, all the request that should
// be sent to all servers, would be sent to one server, which may be a
// disastrous behaviour. In the worst case it would cause the server shutdown
// again if circuit breaker is enabled and the server cluster would never
// recover. This class controls the amount of requests that sent to the revived
// servers when recovering from all servers are shutdown.
class RevivePolicy { class RevivePolicy {
public: public:
// TODO(zhujiashun): // Indicate that reviving from the shutdown of all server is happening.
virtual void StartReviving() = 0;
// Return true if some customized policies are satisfied.
virtual bool DoReject(const std::vector<ServerId>& server_list) = 0;
virtual void StartRevive() = 0; // Stop reviving state and do not reject the request if some condition is
virtual bool RejectDuringReviving(const std::vector<ServerId>& server_list) = 0; // satisfied.
virtual void StopRevivingIfNecessary() = 0;
}; };
// The default revive policy. Once no servers are available, reviving is start.
// If in reviving state, the probability that a request is accepted is q/n, in
// which q is the number of current available server, n is the number of minimum
// working instances setting by user. If q is not changed during a given time,
// hold_time_ms, then the cluster is considered recovered and all the request
// would be sent to the current available servers.
class DefaultRevivePolicy : public RevivePolicy { class DefaultRevivePolicy : public RevivePolicy {
public: public:
DefaultRevivePolicy(int64_t minimum_working_instances, int64_t hold_time_ms); DefaultRevivePolicy(int64_t minimum_working_instances, int64_t hold_time_ms);
void StartRevive() override; void StartReviving();
bool RejectDuringReviving(const std::vector<ServerId>& server_list) override; bool DoReject(const std::vector<ServerId>& server_list);
void StopRevivingIfNecessary();
private: private:
bool _reviving; bool _reviving;
......
...@@ -884,6 +884,7 @@ public: ...@@ -884,6 +884,7 @@ public:
num_reject->fetch_add(1, butil::memory_order_relaxed); num_reject->fetch_add(1, butil::memory_order_relaxed);
} }
} }
delete this;
} }
brpc::Controller cntl; brpc::Controller cntl;
...@@ -898,24 +899,22 @@ TEST_F(LoadBalancerTest, revived_from_all_failed_intergrated) { ...@@ -898,24 +899,22 @@ TEST_F(LoadBalancerTest, revived_from_all_failed_intergrated) {
GFLAGS_NS::SetCommandLineOption("circuit_breaker_short_window_error_percent", "30"); GFLAGS_NS::SetCommandLineOption("circuit_breaker_short_window_error_percent", "30");
GFLAGS_NS::SetCommandLineOption("circuit_breaker_max_isolation_duration_ms", "5000"); GFLAGS_NS::SetCommandLineOption("circuit_breaker_max_isolation_duration_ms", "5000");
const char* lb_algo[] = { "rr" , "random", "wrr", "c_murmurhash" };
char* lb_algo[] = { "rr" , "random" };
brpc::Channel channel; brpc::Channel channel;
brpc::ChannelOptions options; brpc::ChannelOptions options;
options.protocol = "http"; options.protocol = "http";
options.timeout_ms = 300; options.timeout_ms = 300;
options.enable_circuit_breaker = true; options.enable_circuit_breaker = true;
options.revive_policy = new brpc::DefaultRevivePolicy(2, 2000 /*2s*/); options.revive_policy = new brpc::DefaultRevivePolicy(2, 2000 /*2s*/);
// Set max_retry to 0 so that health check of servers // Set max_retry to 0 so that the time of health check of different servers
// are not continuous. // are not continuous.
options.max_retry = 0; options.max_retry = 0;
ASSERT_EQ(channel.Init("list://127.0.0.1:7777,127.0.0.1:7778", ASSERT_EQ(channel.Init("list://127.0.0.1:7777 50,127.0.0.1:7778 50",
lb_algo[butil::fast_rand_less_than(ARRAY_SIZE(lb_algo))], lb_algo[butil::fast_rand_less_than(ARRAY_SIZE(lb_algo))],
&options), 0); &options), 0);
uint64_t request_code = 0;
test::EchoRequest req; test::EchoRequest req;
req.set_message("123"); req.set_message("123");
test::EchoResponse res; test::EchoResponse res;
...@@ -923,12 +922,14 @@ TEST_F(LoadBalancerTest, revived_from_all_failed_intergrated) { ...@@ -923,12 +922,14 @@ TEST_F(LoadBalancerTest, revived_from_all_failed_intergrated) {
// trigger one server to health check // trigger one server to health check
{ {
brpc::Controller cntl; brpc::Controller cntl;
cntl.set_request_code(brpc::policy::MurmurHash32(&++request_code, 8));
stub.Echo(&cntl, &req, &res, NULL); stub.Echo(&cntl, &req, &res, NULL);
} }
bthread_usleep(500000); bthread_usleep(500000);
// trigger the other server to health check // trigger the other server to health check
{ {
brpc::Controller cntl; brpc::Controller cntl;
cntl.set_request_code(brpc::policy::MurmurHash32(&++request_code, 8));
stub.Echo(&cntl, &req, &res, NULL); stub.Echo(&cntl, &req, &res, NULL);
} }
bthread_usleep(500000); bthread_usleep(500000);
...@@ -947,14 +948,18 @@ TEST_F(LoadBalancerTest, revived_from_all_failed_intergrated) { ...@@ -947,14 +948,18 @@ TEST_F(LoadBalancerTest, revived_from_all_failed_intergrated) {
int64_t start_ms = butil::gettimeofday_ms(); int64_t start_ms = butil::gettimeofday_ms();
butil::atomic<int32_t> num_reject(0); butil::atomic<int32_t> num_reject(0);
int64_t q = 0;
while ((butil::gettimeofday_ms() - start_ms) < while ((butil::gettimeofday_ms() - start_ms) <
brpc::FLAGS_health_check_interval * 1000 + 10) { brpc::FLAGS_health_check_interval * 1000 + 10) {
Done* done = new Done; Done* done = new Done;
done->num_reject = &num_reject; done->num_reject = &num_reject;
done->req.set_message("123"); done->req.set_message("123");
done->cntl.set_request_code(brpc::policy::MurmurHash32(&++request_code, 8));
stub.Echo(&done->cntl, &done->req, &done->res, done); stub.Echo(&done->cntl, &done->req, &done->res, done);
q++;
bthread_usleep(1000); bthread_usleep(1000);
} }
ASSERT_TRUE(num_reject.load(butil::memory_order_relaxed) > 1700);
// should recover now // should recover now
butil::atomic<int32_t> num_failed(0); butil::atomic<int32_t> num_failed(0);
...@@ -962,12 +967,12 @@ TEST_F(LoadBalancerTest, revived_from_all_failed_intergrated) { ...@@ -962,12 +967,12 @@ TEST_F(LoadBalancerTest, revived_from_all_failed_intergrated) {
Done* done = new Done; Done* done = new Done;
done->req.set_message("123"); done->req.set_message("123");
done->num_failed = &num_failed; done->num_failed = &num_failed;
done->cntl.set_request_code(brpc::policy::MurmurHash32(&++request_code, 8));
stub.Echo(&done->cntl, &done->req, &done->res, done); stub.Echo(&done->cntl, &done->req, &done->res, done);
bthread_usleep(1000); bthread_usleep(1000);
} }
bthread_usleep(1050*1000 /* sleep longer than timeout of service */); bthread_usleep(1050*1000 /* sleep longer than timeout of service */);
ASSERT_EQ(0, num_failed.load(butil::memory_order_relaxed)); ASSERT_EQ(0, num_failed.load(butil::memory_order_relaxed));
ASSERT_TRUE(num_reject.load(butil::memory_order_relaxed) > 1500);
} }
} //namespace } //namespace
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment