diff --git a/ctl/backup_tar.go b/ctl/backup_tar.go index 0b0780b0d..d0c00b675 100644 --- a/ctl/backup_tar.go +++ b/ctl/backup_tar.go @@ -15,6 +15,7 @@ import ( "path/filepath" "time" + "github.com/djherbis/buffer" pilosa "github.com/featurebasedb/featurebase/v3" "github.com/featurebasedb/featurebase/v3/authn" "github.com/featurebasedb/featurebase/v3/disco" @@ -231,24 +232,7 @@ func (cmd *BackupTarCommand) backupTarIDAllocData(ctx context.Context, tw *tar.W } defer rc.Close() - // Read to buffer to determine size. - if _, err := buf.ReadFrom(rc); err != nil { - return fmt.Errorf("copying id alloc data to memory: %w", err) - } - - // Build header & copy data to archive. - if err = tw.WriteHeader(&tar.Header{ - Name: "idalloc", - Mode: 0o666, - Size: int64(buf.Len()), - ModTime: time.Now(), - }); err != nil { - return err - } else if _, err := io.Copy(tw, buf); err != nil { - return fmt.Errorf("copying id alloc data to archive: %w", err) - } - - return nil + return writeToTar(tw, "idalloc", rc) } // backupTarIndex backs up all shards for a given index. @@ -332,26 +316,7 @@ func (cmd *BackupTarCommand) backupTarShardNode(ctx context.Context, tw *tar.Wri return fmt.Errorf("fetching shard reader: %w", err) } defer rc.Close() - - // Read to buffer to determine size. - // TODO: Provide size via the reader itself. - if _, err := buf.ReadFrom(rc); err != nil { - return fmt.Errorf("copying shard data to memory: %w", err) - } - - // Build header & copy data to archive. - if err = tw.WriteHeader(&tar.Header{ - Name: filename, - Mode: 0o666, - Size: int64(buf.Len()), - ModTime: time.Now(), - }); err != nil { - return err - } else if _, err := io.Copy(tw, buf); err != nil { - return fmt.Errorf("copying shard data to archive: %w", err) - } - - return nil + return writeToTar(tw, filename, rc) } func (cmd *BackupTarCommand) backupTarShardDataframe(ctx context.Context, tw *tar.Writer, indexName string, shard uint64, node *disco.Node, buf *bytes.Buffer) error { @@ -375,23 +340,7 @@ func (cmd *BackupTarCommand) backupTarShardDataframe(ctx context.Context, tw *ta } filename := filepath.Join("indexes", indexName, "dataframe", fmt.Sprintf("%04d", shard)) - logger.Printf("writing %v", filename) - if _, err := buf.ReadFrom(resp.Body); err != nil { - return fmt.Errorf("copying shard data to memory: %w", err) - } - - // Build header & copy data to archive. - if err = tw.WriteHeader(&tar.Header{ - Name: filename, - Mode: 0o666, - Size: int64(buf.Len()), - ModTime: time.Now(), - }); err != nil { - return err - } else if _, err := io.Copy(tw, buf); err != nil { - return fmt.Errorf("copying shard data to archive: %w", err) - } - return nil + return writeToTar(tw, filename, resp.Body) } func (cmd *BackupTarCommand) backupTarIndexTranslateData(ctx context.Context, tw *tar.Writer, name string, buf *bytes.Buffer) error { @@ -418,28 +367,11 @@ func (cmd *BackupTarCommand) backupTarIndexPartitionTranslateData(ctx context.Co } defer rc.Close() - // Read to buffer to determine size. - if _, err := buf.ReadFrom(rc); err != nil { - return fmt.Errorf("copying translate data to memory: %w", err) - } - - // Build header & copy data to archive. - if err = tw.WriteHeader(&tar.Header{ - Name: path.Join("indexes", name, "translate", fmt.Sprintf("%04d", partitionID)), - Mode: 0o666, - Size: int64(buf.Len()), - ModTime: time.Now(), - }); err != nil { - return err - } else if _, err := io.Copy(tw, buf); err != nil { - return fmt.Errorf("copying translate data to archive: %w", err) - } - - return nil + return writeToTar(tw, path.Join("indexes", name, "translate", fmt.Sprintf("%04d", partitionID)), rc) } -func (cmd *BackupTarCommand) backupTarFieldTranslateData(ctx context.Context, tw *tar.Writer, indexName, fieldName string, buf *bytes.Buffer) error { - buf.Reset() +func (cmd *BackupTarCommand) backupTarFieldTranslateData(ctx context.Context, tw *tar.Writer, indexName, fieldName string, buff *bytes.Buffer) error { + // buf.Reset() logger := cmd.Logger() logger.Printf("backing up field translation data: %s/%s", indexName, fieldName) @@ -450,17 +382,37 @@ func (cmd *BackupTarCommand) backupTarFieldTranslateData(ctx context.Context, tw return fmt.Errorf("fetching translate data reader: %w", err) } defer rc.Close() + return writeToTar(tw, path.Join("indexes", indexName, "fields", fieldName, "translate"), rc) +} +func (cmd *BackupTarCommand) TLSHost() string { return cmd.Host } + +func (cmd *BackupTarCommand) TLSConfiguration() server.TLSConfig { return cmd.TLS } + +func writeToTar(tw *tar.Writer, entryName string, rc io.Reader) error { + spillFile, err := os.CreateTemp("", "spill") + if err != nil { + return fmt.Errorf("creating temp file : %w", err) + } + defer func() { + spillFile.Close() + os.Remove(spillFile.Name()) + }() + mb512 := int64(2 << 29) + buf := buffer.NewSpill(buffer.New(mb512), spillFile) + + n, err := io.Copy(buf, rc) // Read to buffer to determine size. - if _, err := buf.ReadFrom(rc); err != nil { + if err != nil { return fmt.Errorf("copying translate data to memory: %w", err) } // Build header & copy data to archive. if err = tw.WriteHeader(&tar.Header{ - Name: path.Join("indexes", indexName, "fields", fieldName, "translate"), - Mode: 0o666, - Size: int64(buf.Len()), + Name: entryName, + Mode: 0o666, + // Size: int64(buf.Len()), + Size: n, ModTime: time.Now(), }); err != nil { return err @@ -469,7 +421,3 @@ func (cmd *BackupTarCommand) backupTarFieldTranslateData(ctx context.Context, tw } return nil } - -func (cmd *BackupTarCommand) TLSHost() string { return cmd.Host } - -func (cmd *BackupTarCommand) TLSConfiguration() server.TLSConfig { return cmd.TLS } diff --git a/go.mod b/go.mod index f92aa51ab..862338d18 100644 --- a/go.mod +++ b/go.mod @@ -83,6 +83,7 @@ require ( github.com/PaesslerAG/gval v1.0.0 github.com/PaesslerAG/jsonpath v0.1.1 github.com/apache/arrow/go/v10 v10.0.0-20221021053532-2f627c213fc3 + github.com/djherbis/buffer v1.2.0 github.com/gomem/gomem v0.1.0 github.com/google/uuid v1.3.0 github.com/jaffee/commandeer v0.6.0 diff --git a/go.sum b/go.sum index fdbd41d61..5f5151ba2 100644 --- a/go.sum +++ b/go.sum @@ -232,6 +232,8 @@ github.com/dgryski/go-farm v0.0.0-20190423205320-6a90982ecee2 h1:tdlZCpZ/P9DhczC github.com/dgryski/go-farm v0.0.0-20190423205320-6a90982ecee2/go.mod h1:SqUrOPUnsFjfmXRMNPybcSiG0BgUW2AuFH8PAnS2iTw= github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f/go.mod h1:cuUVRXasLTGF7a8hSLbxyZXjz+1KgoB3wDUb6vlszIc= github.com/dgryski/go-sip13 v0.0.0-20181026042036-e10d5fee7954/go.mod h1:vAd38F8PWV+bWy6jNmig1y/TA+kYO4g3RSRF0IAv0no= +github.com/djherbis/buffer v1.2.0 h1:PH5Dd2ss0C7CRRhQCZ2u7MssF+No9ide8Ye71nPHcrQ= +github.com/djherbis/buffer v1.2.0/go.mod h1:fjnebbZjCUpPinBRD+TDwXSOeNQ7fPQWLfGQqiAiUyE= github.com/docker/spdystream v0.0.0-20160310174837-449fdfce4d96/go.mod h1:Qh8CwZgvJUkLughtfhJv5dyTYa91l1fOUCrgjqmcifM= github.com/dustin/go-humanize v0.0.0-20171111073723-bb3d318650d4/go.mod h1:HtrtbFcZ19U5GC7JDqmcUSB87Iq5E25KnS6fMYU6eOk= github.com/dustin/go-humanize v1.0.0 h1:VSnTsYCnlFHaM2/igO1h6X3HA71jcobQuxemgkq4zYo=