limit memory usage while calculating entry size

This commit is contained in:
Todd Gruben 2023-02-17 08:58:49 -06:00
parent d0e4012025
commit 0e3af115e7
3 changed files with 34 additions and 83 deletions

View file

@ -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 }

1
go.mod
View file

@ -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

2
go.sum
View file

@ -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=