remove broadcast and use tempfile in stead

This commit is contained in:
Todd Gruben 2023-02-21 11:36:11 -06:00
parent 6253092a2f
commit eac7c49c4b
3 changed files with 25 additions and 10 deletions

View file

@ -23,6 +23,8 @@ The Restore command will take a tar-formatted backup archive and restore it to a
flags.DurationVar(&cmd.RetryPeriod, "retry-period", cmd.RetryPeriod, "Length of time after HTTP request failure to continue retrying request.")
flags.StringVar(&cmd.Pprof, "pprof", cmd.Pprof, "host:port to listen for profiling requests at /debug/pprof and /debug/fgprof.")
flags.StringVar(&cmd.AuthToken, "auth-token", "", "Authentication token")
flags.StringVar(&cmd.TempDir, "temp-dir", cmd.TempDir, "Location of tempory spillover files the default is the system's default(usually /tmp)")
ctl.SetTLSConfig(
flags, "",
&cmd.TLS.CertificatePath,

View file

@ -3,6 +3,7 @@ package ctl
import (
"archive/tar"
"bufio"
"bytes"
"context"
"crypto/tls"
@ -492,3 +493,10 @@ func (fb *FileBuffer) Reset() error {
fb.buf.Reset()
return fb.Close()
}
func (fb *FileBuffer) NewReader() io.Reader {
if fb.file != nil {
return bytes.NewReader(fb.buf.Bytes())
}
return bufio.NewReader(fb.file)
}

View file

@ -3,7 +3,6 @@ package ctl
import (
"archive/tar"
"bytes"
"compress/gzip"
"context"
"crypto/tls"
@ -39,6 +38,8 @@ type RestoreTarCommand struct {
// Host:port on which to listen for pprof.
Pprof string `json:"pprof"`
TempDir string
// Reusable client.
client *pilosa.InternalClient
@ -134,7 +135,10 @@ func (cmd *RestoreTarCommand) Run(ctx context.Context) (err error) {
return errors.New("no primary")
}
c := &gohttp.Client{}
buf := new(bytes.Buffer)
mb512 := 2 << 29
buf := NewFileBuffer(mb512, cmd.TempDir)
defer buf.Close()
for {
buf.Reset()
header, err := tarReader.Next()
@ -188,7 +192,8 @@ func (cmd *RestoreTarCommand) Run(ctx context.Context) (err error) {
node := node
g.Go(func() error {
client := &gohttp.Client{}
rd := bytes.NewReader(buf.Bytes())
// rd := bytes.NewReader(buf.Bytes())
rd := buf.NewReader()
logger.Printf("shard %v %v", shard, indexName)
url := node.URI.Path(fmt.Sprintf("/internal/restore/%v/%v", indexName, shard))
_, err = client.Post(url, "application/octet-stream", rd)
@ -220,7 +225,8 @@ func (cmd *RestoreTarCommand) Run(ctx context.Context) (err error) {
node := node
g.Go(func() error {
client := &gohttp.Client{}
rd := bytes.NewReader(buf.Bytes())
// rd := bytes.NewReader(buf.Bytes())
rd := buf.NewReader()
logger.Printf("dataframe shard %v %v", shard, indexName)
url := node.URI.Path(fmt.Sprintf("/internal/dataframe/restore/%v/%v", indexName, shard))
_, err = client.Post(url, "application/octet-stream", rd)
@ -252,7 +258,8 @@ func (cmd *RestoreTarCommand) Run(ctx context.Context) (err error) {
g.Go(func() error {
// rd := bytes.NewReader(shardBytes)
rd := func() (io.Reader, error) {
return bytes.NewReader(buf.Bytes()), nil
return buf.NewReader(), nil
// return bytes.NewReader(buf.Bytes()), nil
}
return cmd.client.ImportIndexKeys(ctx, &node.URI, indexName, partitionID, false, rd)
@ -269,20 +276,18 @@ func (cmd *RestoreTarCommand) Run(ctx context.Context) (err error) {
switch action := record[4]; action {
case "translate":
logger.Printf("field keys %v %v", indexName, fieldName)
bc := NewBroadcaster(tarReader, len(nodes))
_, err = io.Copy(buf, tarReader)
g, _ := errgroup.WithContext(ctx)
for i, node := range nodes {
i := i
for _, node := range nodes {
node := node
g.Go(func() error {
rd := func() (io.Reader, error) {
return bc.Readers[i], nil
return buf.NewReader(), nil
}
return cmd.client.ImportFieldKeys(ctx, &node.URI, indexName, fieldName, false, rd)
})
}
bc.Consume()
if err := g.Wait(); err != nil {
return err
}