Merge pull request #1677 from nagamocha3000/ha_key_translation

CORE-837 Perform partial replication for index keys
This commit is contained in:
nagamocha3000 2021-09-14 10:53:44 +03:00 • committed by GitHub
commit 248ec8ebff
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
7 changed files with 598 additions and 26 deletions

View file

@ -137,6 +137,13 @@ clustertests: vendor
docker-compose -f $(DOCKER_COMPOSE) up --exit-code-from=client1
DOCKER_COMPOSE_INDEX_KEY_REPLICATION=internal/clustertests/docker-compose-index-key-replication.yml
# Check clustertests target for more info
clustertests-index-key-replication: vendor
docker-compose -f $(DOCKER_COMPOSE_INDEX_KEY_REPLICATION) down
docker-compose -f $(DOCKER_COMPOSE_INDEX_KEY_REPLICATION) build client1
docker-compose -f $(DOCKER_COMPOSE_INDEX_KEY_REPLICATION) up --exit-code-from=client1
# Like clustertests, but rebuilds all images.
clustertests-build: vendor
docker-compose -f $(DOCKER_COMPOSE) down -v

View file

@ -1165,6 +1165,9 @@ func (c *cluster) followResizeInstruction(ctx context.Context, instr *ResizeInst
}
}
// fire off translation sync
_ = c.translationSyncer.Reset()
return nil
}
@ -1179,6 +1182,8 @@ func (c *cluster) resizeAbort() error {
if c.resizeCancel != nil {
c.resizeCancel()
}
// fire off translation sync
_ = c.translationSyncer.Reset()
return nil
}

149
holder.go
View file

@ -1472,7 +1472,10 @@ type holderSyncer struct {
Cluster *cluster
// Translation sync handling.
readers []TranslateEntryReader
readers []TranslateEntryReader
readersMu sync.Mutex
pendingReaders int
stopInitializeReplicationCh chan struct{}
syncers errgroup.Group
@ -1596,6 +1599,23 @@ func (s *holderSyncer) syncFragment(index, field, view string, shard uint64) err
// resetTranslationSync reinitializes streaming sync of translation data.
func (s *holderSyncer) resetTranslationSync() error {
if s.stopInitializeReplicationCh == nil {
// suppose stopTranslationSync[S] holds the lock s.readersMu and tries
// to send a signal to s.stopInitializeRepliationCh. If
// s.stopInitializeReplicationCh is unbufferd and
// s.initializeReplication[I] is trying to acquire s.readersMu, this
// will result in a deadlock.
// Hence, the channel is buffered to prevent this scenario. Once
// [S] releases the lock, [I] will acquire it, however, the value
// of s.pendingReaders will be -1, at this point [I] should not
// attempt to add any more readers since those readers will be 'stale'
// [I] should also drain the channel to prevent the next invocation
// of [I] receiving a stop signal that was meant for the current one
// This is needlessly complicated and is the result of me running
// into various deadlocks while trying to fix handling of column
// key replication.
s.stopInitializeReplicationCh = make(chan struct{})
}
// Stop existing streams.
if err := s.stopTranslationSync(); err != nil {
return errors.Wrap(err, "stop translation sync")
@ -1655,7 +1675,15 @@ func newActiveTranslationSyncer(ch chan struct{}) *activeTranslationSyncer {
// Reset resets the server's translation syncer.
func (a *activeTranslationSyncer) Reset() error {
a.ch <- struct{}{}
// just in case some other part of the code has fired
// off a translation sync and it hasn't been received yet
// therefore we don't want to block on send since a.ch is
// (for now) unbuffered. One translationSync is as good as
// another
select {
case a.ch <- struct{}{}:
default:
}
return nil
}
@ -1665,6 +1693,17 @@ func (a *activeTranslationSyncer) Reset() error {
// to complete. This should be called before reconnecting to the cluster in case
// of a cluster resize or schema change.
func (s *holderSyncer) stopTranslationSync() error {
s.readersMu.Lock()
defer func() {
s.readers = nil // will be populated by initializeReplication
s.readersMu.Unlock()
}()
// send signal to stop initializing more readers
if s.pendingReaders > 0 {
s.pendingReaders = -1
close(s.stopInitializeReplicationCh)
s.stopInitializeReplicationCh = make(chan struct{})
}
var g errgroup.Group
for i := range s.readers {
rd := s.readers[i]
@ -1724,7 +1763,6 @@ func (s *holderSyncer) setTranslateReadOnlyFlags(snap *topology.ClusterSnapshot)
// replicate these from whichever node is primary for that partition).
func (s *holderSyncer) initializeReplication(snap *topology.ClusterSnapshot) error {
nodeMaps := make(map[string]TranslateOffsetMap)
if snap.ReplicaN > 1 {
if err := s.populateIndexReplication(nodeMaps, snap); err != nil {
return err
@ -1734,26 +1772,99 @@ func (s *holderSyncer) initializeReplication(snap *topology.ClusterSnapshot) err
return err
}
// filter out empty nodes
nodes := make(map[*topology.Node]bool)
for _, node := range snap.Nodes {
m := nodeMaps[node.ID]
if m.Empty() {
continue
if !m.Empty() {
nodes[node] = true
}
}
// connect to remote nodes and set up readers
readersCh := make(chan TranslateEntryReader, len(nodes))
ctx, cancelAddingMoreReaders := context.WithCancel(context.Background())
defer cancelAddingMoreReaders()
s.readersMu.Lock()
s.pendingReaders = len(nodes)
s.readersMu.Unlock()
go func() {
for {
for node := range nodes {
// check if ctx cancelled
// this means there was a signal sent to stop further init
// of readers
select {
case <-s.Closing:
return
case <-ctx.Done():
close(readersCh)
return
default:
}
// connect to remote node
m := nodeMaps[node.ID]
rd, err := s.Holder.OpenTranslateReader(context.Background(), node.URI.String(), m)
if err != nil {
continue
}
readersCh <- rd
delete(nodes, node)
}
if len(nodes) == 0 {
close(readersCh)
return
}
time.Sleep(10 * time.Second)
}
}()
for {
select {
case <-s.Closing:
cancelAddingMoreReaders()
return nil
case <-s.stopInitializeReplicationCh:
return nil
case rd, ok := <-readersCh:
// all translate readers have been launched, hence channel is
// closed
if !ok {
return nil
}
s.readersMu.Lock()
// [S] has been initiated and acquired the lock first
// at this point we should close the reader we've recieved rather
// than start replication on it.
// [S] should have already closed all the rest
// of the reads if they were still in action.
// we are also draining the channel since the signal for stopping
// further replication was meant for us
if s.pendingReaders == -1 {
rd.Close()
cancelAddingMoreReaders()
drain:
for {
select {
case <-s.stopInitializeReplicationCh:
default:
break drain
}
}
s.readersMu.Unlock()
return nil
}
s.pendingReaders--
s.readers = append(s.readers, rd)
s.syncers.Go(func() error {
defer rd.Close()
s.readBothTranslateReader(rd, snap)
return nil
})
s.readersMu.Unlock()
}
// Connect to remote node and begin streaming.
rd, err := s.Holder.OpenTranslateReader(context.Background(), node.URI.String(), m)
if err != nil {
return err
}
s.readers = append(s.readers, rd)
s.syncers.Go(func() error {
defer rd.Close()
s.readBothTranslateReader(rd, snap)
return nil
})
}
return nil
}
// populateFieldReplication populates a map from node IDs to TranslateOffsetMaps

View file

@ -21,7 +21,7 @@ import (
"testing"
"time"
"github.com/molecula/featurebase/v2"
pilosa "github.com/molecula/featurebase/v2"
"github.com/molecula/featurebase/v2/disco"
picli "github.com/molecula/featurebase/v2/http"
)

View file

@ -0,0 +1,73 @@
version: '2'
services:
pilosa1:
build:
context: ../..
dockerfile: Dockerfile-clustertests
image: ptest
ports:
- "33455:10101"
environment:
- PILOSA_NAME=pilosa1
- PILOSA_ETCD_LISTEN_CLIENT_ADDRESS=http://0.0.0.0:10201
- PILOSA_ETCD_ADVERTISE_CLIENT_ADDRESS=http://pilosa1:10201
- PILOSA_ETCD_LISTEN_PEER_ADDRESS=http://0.0.0.0:10301
- PILOSA_ETCD_ADVERTISE_PEER_ADDRESS=http://pilosa1:10301
- PILOSA_ETCD_INITIAL_CLUSTER=pilosa1=http://pilosa1:10301,pilosa2=http://pilosa2:10301,pilosa3=http://pilosa3:10301
- PILOSA_CLUSTER_REPLICAS=3
networks:
- pilosanet
command:
- "/featurebase server --bind pilosa1:10101"
pilosa2:
build:
context: ../..
dockerfile: Dockerfile-clustertests
image: ptest
ports:
- "33456:10101"
environment:
- PILOSA_NAME=pilosa2
- PILOSA_ETCD_LISTEN_CLIENT_ADDRESS=http://0.0.0.0:10201
- PILOSA_ETCD_ADVERTISE_CLIENT_ADDRESS=http://pilosa2:10201
- PILOSA_ETCD_LISTEN_PEER_ADDRESS=http://0.0.0.0:10301
- PILOSA_ETCD_ADVERTISE_PEER_ADDRESS=http://pilosa2:10301
- PILOSA_ETCD_INITIAL_CLUSTER=pilosa1=http://pilosa1:10301,pilosa2=http://pilosa2:10301,pilosa3=http://pilosa3:10301
- PILOSA_CLUSTER_REPLICAS=3
networks:
- pilosanet
command:
- "/featurebase server --bind pilosa2:10101"
pilosa3:
build:
context: ../..
dockerfile: Dockerfile-clustertests
image: ptest
ports:
- "33457:10101"
environment:
- PILOSA_NAME=pilosa3
- PILOSA_ETCD_LISTEN_CLIENT_ADDRESS=http://0.0.0.0:10201
- PILOSA_ETCD_ADVERTISE_CLIENT_ADDRESS=http://pilosa3:10201
- PILOSA_ETCD_LISTEN_PEER_ADDRESS=http://0.0.0.0:10301
- PILOSA_ETCD_ADVERTISE_PEER_ADDRESS=http://pilosa3:10301
- PILOSA_ETCD_INITIAL_CLUSTER=pilosa1=http://pilosa1:10301,pilosa2=http://pilosa2:10301,pilosa3=http://pilosa3:10301
- PILOSA_CLUSTER_REPLICAS=3
networks:
- pilosanet
command:
- "/featurebase server --bind pilosa3:10101"
client1:
build:
context: .
environment:
- ENABLE_PILOSA_CLUSTER_TESTS_FOR_INDEX_KEY_REPLICATION=1
- GO111MODULE=on
networks:
- pilosanet
volumes:
- /var/run/docker.sock:/var/run/docker.sock
command:
- "cd /go/src/github.com/molecula/featurebase/ && go test -mod=vendor -v -run=IndexKey -count=1 github.com/molecula/featurebase/v2/internal/clustertests"
networks:
pilosanet:

View file

@ -0,0 +1,375 @@
// Copyright 2017 Pilosa Corp.
//
// 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.
package clustertest
import (
"context"
"crypto/tls"
"fmt"
"math/rand"
"os"
"os/exec"
"sync"
"testing"
"time"
pilosa "github.com/molecula/featurebase/v2"
"github.com/molecula/featurebase/v2/disco"
"github.com/molecula/featurebase/v2/http"
picli "github.com/molecula/featurebase/v2/http"
"github.com/molecula/featurebase/v2/net"
"github.com/molecula/featurebase/v2/topology"
)
// index -> key -> ids from all replicas
type translationRes map[string]map[string][]uint64
func defaultTranslationResults(indexes []string) translationRes {
res := make(translationRes)
for _, index := range indexes {
res[index] = make(map[string][]uint64)
}
return res
}
func verify(allRes translationRes, indexes []string, replicasN, count int) error {
received := 0
allIDsSame := func(ids []uint64) bool {
if len(ids) <= 1 {
// trivially true
return true
}
id := ids[0]
for _, other := range ids {
if id != other {
return false
}
}
return true
}
for _, index := range indexes {
res, ok := allRes[index]
if !ok {
return fmt.Errorf("Expected index '%s' but not present", index)
}
for key, ids := range res {
// first id is after successful translation, rest are from
// translate readers
replicationCount := len(ids) - 1
received += replicationCount
if replicationCount != replicasN {
return fmt.Errorf("Count for ids for key '%s'(%d), index '%s' not equal than replicasN(%d)", key, replicationCount, index, replicasN)
}
if !allIDsSame(ids) {
return fmt.Errorf("Expected all ids for key '%s', index '%s' to be the same across the cluster %+v", key, index, ids)
}
}
}
if received != count {
return fmt.Errorf("Expected %d count of keys, received %d", count, received)
}
return nil
}
func getURIsFromAddresses(addrs []string) ([]*net.URI, error) {
uris := make([]*net.URI, 0, len(addrs))
for _, addr := range addrs {
uri, err := net.NewURIFromAddress(addr)
if err != nil {
return nil, err
}
uris = append(uris, uri)
}
return uris, nil
}
func getClients(addrs []string) ([]*http.InternalClient, error) {
clients := make([]*http.InternalClient, 0, len(addrs))
for _, addr := range addrs {
c, err := picli.NewInternalClient(addr, picli.GetHTTPClient(nil))
if err != nil {
return nil, err
}
clients = append(clients, c)
}
return clients, nil
}
func genIndexNames(indexCount int) []string {
indexNames := make([]string, 0, indexCount)
for i := 1; i <= indexCount; i++ {
indexNames = append(indexNames, fmt.Sprintf("idx-%d", i))
}
return indexNames
}
func parseDuration(t *testing.T, s string) time.Duration {
parsed, err := time.ParseDuration(s)
if err != nil {
t.Fatal(err)
}
return parsed
}
func durationToSeconds(d time.Duration) string {
s := int(d.Seconds())
return fmt.Sprintf("%ds", s)
}
func TestIndexKeyReplication(t *testing.T) {
if os.Getenv("ENABLE_PILOSA_CLUSTER_TESTS_FOR_INDEX_KEY_REPLICATION") != "1" {
t.Skip("pilosa cluster tests for index key replication are not enabled")
}
// configurations for test
replicasN := 3
indexCount := 4
intervalDurationArg := "100ms"
totalInsertionDurationArg := "10s"
coolOffDurationArg := "5s"
numKeysToInsertPerDuration := 100
addresses := []string{"pilosa1:10101", "pilosa2:10101", "pilosa3:10101"}
intervalDuration := parseDuration(t, intervalDurationArg)
totalInsertionDuration := parseDuration(t, totalInsertionDurationArg)
coolOffDuration := parseDuration(t, coolOffDurationArg)
indexes := genIndexNames(indexCount)
clients, err := getClients(addresses)
cli := clients[0]
if err != nil {
t.Fatalf("on init clients from addresses: %v, %v", addresses, err)
}
uris, err := getURIsFromAddresses(addresses)
if err != nil {
t.Fatalf("on init clients from addresses: %v, %v", addresses, err)
}
ctx := context.Background()
// create index keyed
for _, index := range indexes {
err = cli.EnsureIndex(ctx, index, pilosa.IndexOptions{
Keys: true,
})
if err != nil {
t.Fatalf("creating/asserting index: %v", err)
}
}
// set up index keys to insert
keysInserted := 0
allRes := defaultTranslationResults(indexes)
{
ctxForIndexKeyCreation, cancelFurtherInsertions := context.WithCancel(ctx)
var wg sync.WaitGroup
var mu = &sync.Mutex{}
wg.Add(len(indexes))
seed := time.Now().UnixNano()
t.Logf("start inserting index keys for %v, seed(%v)", totalInsertionDuration, seed)
rng := rand.New(rand.NewSource(seed))
for _, index := range indexes {
translations := allRes[index]
go func(index string, translations map[string][]uint64) {
keys := make([]string, numKeysToInsertPerDuration)
offset := 1
ticker := time.NewTicker(intervalDuration)
defer func() {
ticker.Stop()
wg.Done()
}()
for {
select {
case <-ticker.C:
for i := 0; i < numKeysToInsertPerDuration; i++ {
keys[i] = fmt.Sprintf("key-%d", i+offset)
}
// pick random node to send key creation to
r := rng.Intn(len(addresses))
cli, uri := clients[r], uris[r]
// insert index keys
transmap, err := cli.CreateIndexKeysNode(ctxForIndexKeyCreation, uri, index, keys...)
if err != nil {
if err == context.Canceled {
return
} else {
t.Logf("creating index keys for index(%s) send to node(%s): %v", index, uri.String(), err)
continue
}
}
for key, id := range transmap {
translations[key] = append(translations[key], id)
}
mu.Lock()
keysInserted += len(transmap)
mu.Unlock()
offset += numKeysToInsertPerDuration
case <-ctxForIndexKeyCreation.Done():
return
}
}
}(index, translations)
}
// inject fault
pcmd := exec.Command("/pumba", "netem", "--duration",
durationToSeconds(totalInsertionDuration),
"loss", "--percent", "50", "--correlation", "60",
"clustertests_pilosa3_1")
pcmd.Stdout = os.Stdout
pcmd.Stderr = os.Stderr
t.Logf("sending pumba fault injection cmd: %v", pcmd.String())
err = pcmd.Start()
if err != nil {
t.Fatalf("starting pumba command: %v", err)
}
err = pcmd.Wait()
if err != nil {
t.Fatalf("waiting on pumba pause cmd: %v", err)
}
// wait for index keys to be created
t.Logf("start wait to complete index key creation")
time.Sleep(totalInsertionDuration)
cancelFurtherInsertions()
wg.Wait()
t.Logf("done with inserting index keys. Total keys inserted: %d", keysInserted)
}
t.Logf("start cool off period: %v\n", coolOffDuration)
time.Sleep(coolOffDuration)
t.Log("done with cool off period, waiting for stability")
waitForStatus(t, clients[0].Status, string(disco.ClusterStateNormal), 30, 1*time.Second)
t.Log("done with waiting for stability, starting verifying persistence of index keys")
// get all nodes
nodes, err := cli.Nodes(ctx)
if err != nil {
t.Fatal(err)
}
// prepare translate offset maps
nodeMaps := make(map[string]pilosa.TranslateOffsetMap)
for _, n := range nodes {
nodeMaps[n.ID] = make(pilosa.TranslateOffsetMap)
}
schema, err := cli.Schema(ctx)
if err != nil {
t.Fatal(err)
}
for _, indexInfo := range schema {
index := indexInfo.Name
isKeyed := indexInfo.Options.Keys
if !isKeyed {
continue
}
partitionN := topology.DefaultPartitionN
for partition := 0; partition < partitionN; partition++ {
nodes, err := cli.PartitionNodes(ctx, partition)
if err != nil {
t.Fatal(err)
}
for _, n := range nodes {
m := nodeMaps[n.ID]
m.SetIndexPartitionOffset(index, partition, 0)
}
}
}
// open translate reader
readers := make([]pilosa.TranslateEntryReader, 0, len(nodes))
closeAllTranslateReaders := func() []error {
var errs []error
for _, tr := range readers {
err := tr.Close()
if err != nil {
errs = append(errs, err)
}
}
return errs
}
var tlsConfig *tls.Config = nil
var wg sync.WaitGroup
entries := make(chan *pilosa.TranslateEntry)
for _, n := range nodes {
client := http.GetHTTPClient(tlsConfig)
nodeURL := n.URI.String()
offsets := nodeMaps[n.ID]
openTranslateReader := http.GetOpenTranslateReaderFunc(client)
tr, err := openTranslateReader(ctx, nodeURL, offsets)
if err != nil {
t.Fatal(err)
}
wg.Add(1)
readers = append(readers, tr)
go func(node *topology.Node, tr pilosa.TranslateEntryReader) {
defer func() {
wg.Done()
}()
for {
var entry pilosa.TranslateEntry
err := tr.ReadEntry(&entry)
if err != nil {
// TODO also ignore transient http: read on clsoed response body errors
// for now just print error
if err != context.Canceled {
fmt.Printf("node(%s) On read from translate entry reader: %v", node.URI.String(), err)
}
return
}
// ignore field keys
if entry.Field != "" {
continue
}
entries <- &entry
}
}(n, tr)
}
// receive all translation entries made
count := replicasN * keysInserted
i := 0
for {
entry := <-entries
res, ok := allRes[entry.Index]
if !ok {
// ignore indexes we did not create for this test
continue
}
ids, ok := res[entry.Key]
if !ok {
// ignore keys we did not insert for the test
continue
}
res[entry.Key] = append(ids, entry.ID)
i++
if i == count {
break
}
}
errs := closeAllTranslateReaders()
wg.Wait()
for _, err := range errs {
if err != nil {
t.Errorf("on close translate readers: %v", err)
}
}
close(entries)
err = verify(allRes, indexes, replicasN, count)
if err != nil {
t.Fatal(err)
}
}

View file

@ -437,7 +437,7 @@ func NewServer(opts ...ServerOption) (*Server, error) {
confirmDownRetries: defaultConfirmDownRetries,
confirmDownSleep: defaultConfirmDownSleep,
resetTranslationSyncCh: make(chan struct{}),
resetTranslationSyncCh: make(chan struct{}, 1),
logger: logger.NopLogger,
}
@ -569,11 +569,6 @@ func (s *Server) Open() error {
log.Println(errors.Wrap(err, "logging startup"))
}
// Start background process listening for translation
// sync resets.
s.wg.Add(1)
go func() { defer s.wg.Done(); s.monitorResetTranslationSync() }()
// Start DisCo.
ctx, cancel := context.WithTimeout(context.Background(), 120*time.Second)
defer cancel()
@ -612,6 +607,12 @@ func (s *Server) Open() error {
s.syncer.Closing = s.closing
s.syncer.Stats = s.holder.Stats.WithTags("component:HolderSyncer")
// Start background process listening for translation
// sync resets.
s.wg.Add(1)
go func() { defer s.wg.Done(); s.monitorResetTranslationSync() }()
go func() { _ = s.translationSyncer.Reset() }()
// Open holder.
func() {
s.holder.startMsgsMu.Lock()