diff --git a/admin_client.go b/admin_client.go index 8dbda26d..72c04d4c 100644 --- a/admin_client.go +++ b/admin_client.go @@ -68,6 +68,11 @@ func newAdminClient(zkquorum string, options ...Option) AdminClient { option(c) } c.zkClient = zk.NewClient(zkquorum, c.zkTimeout) + + // Get client connection for admin region + c.adminRegionInfo.MarkUnavailable() + go c.reestablishRegion(c.adminRegionInfo) + return c } diff --git a/caches.go b/caches.go index 2d44042f..c367f1dd 100644 --- a/caches.go +++ b/caches.go @@ -62,7 +62,6 @@ func (rcc *clientRegionCache) del(r hrpc.RegionInfo) { rcc.m.Lock() c := r.Client() if c != nil { - r.SetClient(nil) regions := rcc.regions[c] delete(regions, r) } @@ -71,11 +70,7 @@ func (rcc *clientRegionCache) del(r hrpc.RegionInfo) { func (rcc *clientRegionCache) closeAll() { rcc.m.Lock() - for client, regions := range rcc.regions { - for region := range regions { - region.MarkUnavailable() - region.SetClient(nil) - } + for client := range rcc.regions { client.Close() } rcc.m.Unlock() diff --git a/client.go b/client.go index e370234e..404c0ddb 100644 --- a/client.go +++ b/client.go @@ -142,6 +142,10 @@ func newClient(zkquorum string, options ...Option) *client { //since the zkTimeout could be changed as an option c.zkClient = zk.NewClient(zkquorum, c.zkTimeout) + // Get client connection for meta region + c.metaRegionInfo.MarkUnavailable() + go c.reestablishRegion(c.metaRegionInfo) + return c } diff --git a/debug_state_test.go b/debug_state_test.go index fe6bcd23..5886a657 100644 --- a/debug_state_test.go +++ b/debug_state_test.go @@ -44,7 +44,7 @@ func TestDebugStateSanity(t *testing.T) { } else if len(os) != 0 { t.Errorf("Didn't expect any overlaps, got: %v", os) } - region1.SetClient(regClient) + setRegionClient(region1, regClient) client.clients.put("regionserver:1", region1, newClientFn) region2 := region.NewInfo( @@ -60,7 +60,7 @@ func TestDebugStateSanity(t *testing.T) { } else if len(os) != 0 { t.Errorf("Didn't expect any overlaps, got: %v", os) } - region2.SetClient(regClient) + setRegionClient(region2, regClient) client.clients.put("regionserver:1", region2, newClientFn) region3 := region.NewInfo( @@ -76,7 +76,7 @@ func TestDebugStateSanity(t *testing.T) { } else if len(os) != 0 { t.Errorf("Didn't expect any overlaps, got: %v", os) } - region3.SetClient(regClient) + setRegionClient(region3, regClient) client.clients.put("regionserver:1", region3, newClientFn) jsonVal, err := DebugState(client) diff --git a/hrpc/call.go b/hrpc/call.go index 8e84d56e..ea9e36f2 100644 --- a/hrpc/call.go +++ b/hrpc/call.go @@ -17,11 +17,29 @@ import ( ) // RegionInfo represents HBase region. +// +// A RegionInfo may be in one of three states: +// +// 1. Unavailable: AvailabilityChan() is non-nil and Client() is nil +// 2. Available: AvailabilityChan() is nil and Client() is non-nil +// 3. Stale: AvailabilityChan() is nil and Client() is nil +// +// RegionInfo lifecycle: +// +// On creation and when seeing certain errors a region should have +// MarkUnavailable called to initialize the AvailabilityChan. If +// MarkUnavailable returns true then (re)establishRegion will need to +// be called to find a client for the region. The user of the +// RegionInfo can watch the AvailabilityChan to know when the region +// has been established. After the AvailabilityChan is closed Client() +// may return an initialized RegionClient or it may return nil, in +// which case the region is stale and the region lookup should be +// performed again. type RegionInfo interface { IsUnavailable() bool AvailabilityChan() <-chan struct{} MarkUnavailable() bool - MarkAvailable() + MarkAvailable(RegionClient) MarkDead() Context() context.Context String() string @@ -31,7 +49,6 @@ type RegionInfo interface { StopKey() []byte Namespace() []byte Table() []byte - SetClient(RegionClient) Client() RegionClient } diff --git a/hrpc/hrpc_test.go b/hrpc/hrpc_test.go index 12a0f5df..b0069ee2 100644 --- a/hrpc/hrpc_test.go +++ b/hrpc/hrpc_test.go @@ -482,7 +482,7 @@ func (ri mockRegionInfo) Name() []byte { func (ri mockRegionInfo) IsUnavailable() bool { return true } func (ri mockRegionInfo) AvailabilityChan() <-chan struct{} { return nil } func (ri mockRegionInfo) MarkUnavailable() bool { return true } -func (ri mockRegionInfo) MarkAvailable() {} +func (ri mockRegionInfo) MarkAvailable(RegionClient) {} func (ri mockRegionInfo) MarkDead() {} func (ri mockRegionInfo) Context() context.Context { return nil } func (ri mockRegionInfo) String() string { return "" } @@ -491,7 +491,6 @@ func (ri mockRegionInfo) StartKey() []byte { return nil } func (ri mockRegionInfo) StopKey() []byte { return nil } func (ri mockRegionInfo) Namespace() []byte { return nil } func (ri mockRegionInfo) Table() []byte { return nil } -func (ri mockRegionInfo) SetClient(RegionClient) {} func (ri mockRegionInfo) Client() RegionClient { return nil } type byFamily []*pb.MutationProto_ColumnValue diff --git a/region/info.go b/region/info.go index 2b6a3431..729a3066 100644 --- a/region/info.go +++ b/region/info.go @@ -177,6 +177,7 @@ func (i *info) MarkUnavailable() bool { created := false i.m.Lock() if i.available == nil { + i.client = nil i.available = make(chan struct{}) created = true } @@ -186,10 +187,11 @@ func (i *info) MarkUnavailable() bool { // MarkAvailable will mark this region as available again, by closing the struct // returned by AvailabilityChan -func (i *info) MarkAvailable() { +func (i *info) MarkAvailable(c hrpc.RegionClient) { i.m.Lock() ch := i.available i.available = nil + i.client = c close(ch) i.m.Unlock() } @@ -254,7 +256,7 @@ func (i *info) Client() hrpc.RegionClient { return c } -// SetClient sets region client +// SetClient sets region client. Used in tests only func (i *info) SetClient(c hrpc.RegionClient) { i.m.Lock() i.client = c diff --git a/rpc.go b/rpc.go index ac6f426f..d734e347 100644 --- a/rpc.go +++ b/rpc.go @@ -128,31 +128,9 @@ func (c *client) getRegionAndClientForRPC(ctx context.Context, rpc hrpc.Call) ( client := reg.Client() if client == nil { - // There was an error getting the region client. Mark the - // region as unavailable. - if reg.MarkUnavailable() { - // If this was the first goroutine to mark the region as - // unavailable, start a goroutine to reestablish a connection - go c.reestablishRegion(reg) - } - if ch := reg.AvailabilityChan(); ch != nil { - select { - case <-ctx.Done(): - return nil, ctx.Err() - case <-c.done: - return nil, ErrClientClosed - case <-ch: - } - } - if reg.Context().Err() != nil { - // region is dead because it was split or merged, - // retry lookup - continue - } - client = reg.Client() - if client == nil { - continue - } + // region client is nil after marked available, likely + // means the region is dead. Let's retry + continue } rpc.SetRegion(reg) return client, nil @@ -399,7 +377,6 @@ func (c *client) clientDown(client hrpc.RegionClient) { downregions := c.clients.clientDown(client) for downreg := range downregions { if downreg.MarkUnavailable() { - downreg.SetClient(nil) go c.reestablishRegion(downreg) } } @@ -656,6 +633,20 @@ func isRegionEstablished(rc hrpc.RegionClient, reg hrpc.RegionInfo) error { } } +// establishRegion will attempt to find a region client for reg and +// create a connection to that region client. +// +// reg must have had MarkUnavailable called on it before being passed +// to establishRegion to have its AvailabilityChan created. +// establishRegion will call MarkAvailable on the region before +// returning, closing the AvailabilityChan. +// +// When successfully connected to a region server, the region will +// have its Client set. On error it will retry. If the region is found +// to be stale, eg. it has been replaced with a newer region due to a +// split, or otherwise is not found in the meta region then the region +// will have a nil Client associated with it. Callers will need to +// retry lookup. func (c *client) establishRegion(reg hrpc.RegionInfo, addr string) { var backoff time.Duration var err error @@ -663,7 +654,7 @@ func (c *client) establishRegion(reg hrpc.RegionInfo, addr string) { backoff, err = sleepAndIncreaseBackoff(reg.Context(), backoff) if err != nil { // region is dead - reg.MarkAvailable() + reg.MarkAvailable(nil) // unblock waiters return } if addr == "" { @@ -677,7 +668,7 @@ func (c *client) establishRegion(reg hrpc.RegionInfo, addr string) { // region doesn't exist, delete it from caches c.regions.del(originalReg) c.clients.del(originalReg) - originalReg.MarkAvailable() + originalReg.MarkAvailable(nil) // unblock waiters log.WithFields(log.Fields{ "region": originalReg.String(), @@ -686,9 +677,10 @@ func (c *client) establishRegion(reg hrpc.RegionInfo, addr string) { }).Info("region does not exist anymore") return - } else if originalReg.Context().Err() != nil { + } + if originalReg.Context().Err() != nil { // region is dead - originalReg.MarkAvailable() + originalReg.MarkAvailable(nil) // unblock waiters log.WithFields(log.Fields{ "region": originalReg.String(), @@ -697,10 +689,12 @@ func (c *client) establishRegion(reg hrpc.RegionInfo, addr string) { }).Info("region became dead while establishing client for it") return - } else if err == ErrClientClosed { - // client has been closed + } + if err == ErrClientClosed { + // gohbase client has been closed return - } else if err != nil { + } + if err != nil { log.WithFields(log.Fields{ "region": originalReg.String(), "err": err, @@ -714,8 +708,8 @@ func (c *client) establishRegion(reg hrpc.RegionInfo, addr string) { overlaps, replaced := c.regions.put(reg) if !replaced { // a region that is the same or younger is already in cache - reg.MarkAvailable() - originalReg.MarkAvailable() + reg.MarkAvailable(nil) + originalReg.MarkAvailable(nil) return } // otherwise delete the overlapped regions in cache @@ -724,7 +718,7 @@ func (c *client) establishRegion(reg hrpc.RegionInfo, addr string) { } // let rpcs know that they can retry and either get the newly // added region from cache or lookup the one they need - originalReg.MarkAvailable() + originalReg.MarkAvailable(nil) } else { // same region, discard the looked up one reg = originalReg @@ -754,16 +748,12 @@ func (c *client) establishRegion(reg hrpc.RegionInfo, addr string) { if err == nil { if reg == c.adminRegionInfo { - reg.SetClient(client) - reg.MarkAvailable() + reg.MarkAvailable(client) return } if err = isRegionEstablished(client, reg); err == nil { - // set region client so that as soon as we mark it available, - // concurrent readers are able to find the client - reg.SetClient(client) - reg.MarkAvailable() + reg.MarkAvailable(client) return } else if _, ok := err.(region.ServerError); ok { // the client we got died @@ -771,7 +761,7 @@ func (c *client) establishRegion(reg hrpc.RegionInfo, addr string) { } } else if err == context.Canceled { // region is dead - reg.MarkAvailable() + reg.MarkAvailable(nil) return } else { // otherwise Dial failed, purge the client and retry. diff --git a/rpc_test.go b/rpc_test.go index f4b1dd80..01d45173 100644 --- a/rpc_test.go +++ b/rpc_test.go @@ -42,7 +42,7 @@ func newRegionClientFn(addr string) func() hrpc.RegionClient { } func newMockClient(zkClient zk.Client) *client { - return &client{ + c := &client{ clientType: region.RegionClient, regions: keyRegionCache{regions: b.TreeNew[[]byte, hrpc.RegionInfo](region.Compare)}, clients: clientRegionCache{ @@ -58,6 +58,12 @@ func newMockClient(zkClient zk.Client) *client { regionReadTimeout: region.DefaultReadTimeout, newRegionClientFn: newMockRegionClient, } + + return c +} + +func setRegionClient(reg hrpc.RegionInfo, rc hrpc.RegionClient) { + reg.(interface{ SetClient(c hrpc.RegionClient) }).SetClient(rc) } func TestSendRPCSanity(t *testing.T) { @@ -66,7 +72,11 @@ func TestSendRPCSanity(t *testing.T) { // we expect to ask zookeeper for where metaregion is zkClient := mockZk.NewMockClient(ctrl) zkClient.EXPECT().LocateResource(zk.Meta).Return("regionserver:1", nil).MinTimes(1) + c := newMockClient(zkClient) + // Establish meta region like newClient would + c.metaRegionInfo.MarkUnavailable() + go c.reestablishRegion(c.metaRegionInfo) // ask for "theKey" in table "test" mockCall := mock.NewMockCall(ctrl) @@ -151,14 +161,14 @@ func TestReestablishRegionSplit(t *testing.T) { // marking unavailable to simulate error origlReg.MarkUnavailable() rc1 := c.clients.put("regionserver:1", origlReg, newRegionClientFn("regionserver:1")) - origlReg.SetClient(rc1) + setRegionClient(origlReg, rc1) c.regions.put(origlReg) rc2 := c.clients.put("regionserver:1", c.metaRegionInfo, newRegionClientFn("regionserver:1")) if rc1 != rc2 { t.Fatal("expected to get the same region client") } - c.metaRegionInfo.SetClient(rc2) + setRegionClient(c.metaRegionInfo, rc2) c.reestablishRegion(origlReg) @@ -244,8 +254,8 @@ func TestReestablishRegionNSRE(t *testing.T) { } // "nsre" is at the moment at regionserver:1 - c.metaRegionInfo.SetClient(rc1) - origlReg.SetClient(rc1) + setRegionClient(c.metaRegionInfo, rc1) + setRegionClient(origlReg, rc1) // marking unavailable to simulate error origlReg.MarkUnavailable() c.regions.put(origlReg) @@ -321,7 +331,7 @@ func TestEstablishRegionDialFail(t *testing.T) { // inject a fake regionserver client and fake region into cache // pretend regionserver:0 has meta table rc1 := c.clients.put("regionserver:0", c.metaRegionInfo, newRegionClientFn("regionserver:0")) - c.metaRegionInfo.SetClient(rc1) + setRegionClient(c.metaRegionInfo, rc1) // should get stuck if the region is never established c.establishRegion(reg, "regionserver:1") @@ -395,7 +405,7 @@ func TestEstablishServerErrorDuringProbe(t *testing.T) { // pretend regionserver:0 has meta table rc := c.clients.put("regionserver:0", c.metaRegionInfo, newRegionClientFn("regionserver:0")) - c.metaRegionInfo.SetClient(rc) + setRegionClient(c.metaRegionInfo, rc) mockCall := mock.NewMockCall(ctrl) mockCall.EXPECT().Context().Return(context.Background()).AnyTimes() @@ -443,7 +453,7 @@ func TestSendRPCToRegionClientDownDelayed(t *testing.T) { c.clients.put("regionserver:0", origlReg, func() hrpc.RegionClient { return rc }) - origlReg.SetClient(rc) + setRegionClient(origlReg, rc) mockCall := mock.NewMockCall(ctrl) mockCall.EXPECT().Region().Return(origlReg).Times(1) @@ -461,7 +471,7 @@ func TestSendRPCToRegionClientDownDelayed(t *testing.T) { c.clients.put("regionserver:0", origlReg, func() hrpc.RegionClient { return rc2 }) - origlReg.SetClient(rc2) + setRegionClient(origlReg, rc2) // return ServerError from QueueRPC, to emulate dead client result <- hrpc.RPCResult{Error: region.ServerError{}} @@ -504,7 +514,7 @@ func TestReestablishDeadRegion(t *testing.T) { c.clients.put("regionserver:0", reg, newRegionClientFn("regionserver:0")) // pretend regionserver:0 has meta table - c.metaRegionInfo.SetClient(rc1) + setRegionClient(c.metaRegionInfo, rc1) reg.MarkUnavailable() @@ -584,7 +594,7 @@ func TestFindRegion(t *testing.T) { c := newMockClient(nil) // pretend regionserver:0 has meta table rc := c.clients.put("regionserver:0", c.metaRegionInfo, newRegionClientFn("regionserver:0")) - c.metaRegionInfo.SetClient(rc) + setRegionClient(c.metaRegionInfo, rc) ctx := context.Background() testTable := []byte("test") @@ -645,14 +655,14 @@ func TestErrCannotFindRegion(t *testing.T) { // pretend regionserver:0 has meta table rc := c.clients.put("regionserver:0", c.metaRegionInfo, newRegionClientFn("regionserver:0")) - c.metaRegionInfo.SetClient(rc) + setRegionClient(c.metaRegionInfo, rc) // add young and small region to cache origlReg := region.NewInfo(1434573235910, nil, []byte("test"), []byte("test,yolo,1434573235910.56f833d5569a27c7a43fbf547b4924a4."), []byte("yolo"), nil) c.regions.put(origlReg) rc = c.clients.put("regionserver:0", origlReg, newRegionClientFn("regionserver:0")) - origlReg.SetClient(rc) + setRegionClient(origlReg, rc) // request a key not in the "yolo" region. get, err := hrpc.NewGetStr(context.Background(), "test", "meow") @@ -672,7 +682,7 @@ func TestMetaLookupTableNotFound(t *testing.T) { c := newMockClient(nil) // pretend regionserver:0 has meta table rc := c.clients.put("regionserver:0", c.metaRegionInfo, newRegionClientFn("regionserver:0")) - c.metaRegionInfo.SetClient(rc) + setRegionClient(c.metaRegionInfo, rc) _, _, err := c.metaLookup(context.Background(), []byte("tablenotfound"), []byte(t.Name())) if err != TableNotFound { @@ -684,7 +694,7 @@ func TestMetaLookupCanceledContext(t *testing.T) { c := newMockClient(nil) // pretend regionserver:0 has meta table rc := c.clients.put("regionserver:0", c.metaRegionInfo, newRegionClientFn("regionserver:0")) - c.metaRegionInfo.SetClient(rc) + setRegionClient(c.metaRegionInfo, rc) ctx, cancel := context.WithCancel(context.Background()) cancel() @@ -702,6 +712,10 @@ func TestConcurrentRetryableError(t *testing.T) { // keep failing on zookeeper lookup zkc.EXPECT().LocateResource(gomock.Any()).Return("", errors.New("ooops")).AnyTimes() c := newMockClient(zkc) + // Establish meta region like newClient would + c.metaRegionInfo.MarkUnavailable() + go c.reestablishRegion(c.metaRegionInfo) + // create region with mock clien origlReg := region.NewInfo( 0, @@ -731,8 +745,8 @@ func TestConcurrentRetryableError(t *testing.T) { c.regions.put(whateverRegion) c.clients.put("host:1234", origlReg, newRC) c.clients.put("host:1234", whateverRegion, newRC) - origlReg.SetClient(rc) - whateverRegion.SetClient(rc) + setRegionClient(origlReg, rc) + setRegionClient(whateverRegion, rc) numCalls := 100 rc.EXPECT().QueueRPC(gomock.Any()).MinTimes(1) @@ -830,6 +844,9 @@ func TestSendBatchBasic(t *testing.T) { zkClient := mockZk.NewMockClient(ctrl) zkClient.EXPECT().LocateResource(zk.Meta).Return("regionserver:1", nil).MinTimes(1) c := newMockClient(zkClient) + // Establish meta region like newClient would + c.metaRegionInfo.MarkUnavailable() + go c.reestablishRegion(c.metaRegionInfo) call, err := hrpc.NewPutStr(context.Background(), "test", "theKey", nil) if err != nil { @@ -893,7 +910,6 @@ func TestSendBatchBadInput(t *testing.T) { defer ctrl.Finish() zkc := mockZk.NewMockClient(ctrl) - zkc.EXPECT().LocateResource(zk.Meta).Return("regionserver:1", nil).AnyTimes() c := newMockClient(zkc) newRPC := func(table string, batchable bool) hrpc.Call { @@ -984,11 +1000,11 @@ func TestFindClients(t *testing.T) { c := newMockClient(nil) // pretend regionserver:0 has meta table rc := c.clients.put("regionserver:0", c.metaRegionInfo, newRegionClientFn("regionserver:0")) - c.metaRegionInfo.SetClient(rc) + setRegionClient(c.metaRegionInfo, rc) registerRegion := func(reg hrpc.RegionInfo, addr string) { rc := c.clients.put(addr, reg, newRegionClientFn(addr)) - reg.SetClient(rc) + setRegionClient(reg, rc) overlaps, replaced := c.regions.put(reg) if len(overlaps) > 0 { t.Fatalf("overlaps: %v replaced: %t", overlaps, replaced) @@ -1179,11 +1195,11 @@ func TestSendBatchWaitForCompletion(t *testing.T) { c := newMockClient(nil) // pretend regionserver:0 has meta table rc := c.clients.put("regionserver:0", c.metaRegionInfo, newRegionClientFn("regionserver:0")) - c.metaRegionInfo.SetClient(rc) + setRegionClient(c.metaRegionInfo, rc) registerRegion := func(reg hrpc.RegionInfo, addr string) { rc := c.clients.put(addr, reg, newRegionClientFn(addr)) - reg.SetClient(rc) + setRegionClient(reg, rc) overlaps, replaced := c.regions.put(reg) if len(overlaps) > 0 { t.Fatalf("overlaps: %v replaced: %t", overlaps, replaced)