Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 4 additions & 3 deletions api/api.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,9 +36,10 @@ type VersionResponse struct {
}

type PoolPostRequest struct {
Name string `json:"name"`
SortKeys SortKeys `json:"layout"`
Thresh int64 `json:"thresh"`
Name string `json:"name"`
SortKeys SortKeys `json:"layout"`
ObjectCap uint64 `json:"objectcap"`
FrameCap uint64 `json:"framecap"`
Comment on lines +41 to +42

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: I don't feel strongly about this but knobs like these are usually suffixed with "lim", "limit", or "max" so one of those instead of "cap" might make their behavior a little clearer to system users and code readers. (My initial reaction was, "'Cap' probably means maximum here but if that were the case it'd just be 'max" so maybe it means something else.")

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Second nit:

Suggested change
ObjectCap uint64 `json:"objectcap"`
FrameCap uint64 `json:"framecap"`
ObjectCap uint64 `json:"object_cap"`
FrameCap uint64 `json:"frame_cap"`

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm gonna leave these as is since we will discuss and replace them all in a subsequent PR. No since replacing all the tests now and changing again.

}

type SortKeys struct {
Expand Down
53 changes: 31 additions & 22 deletions bsup/writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,45 +14,52 @@ import (
"github.com/superdb/super/vector/vio"
)

var maxFrameSize uint32 = 120_000

// XXX a future PR will wire in compress / thresh options to command line.
// XXX Rows is the key flag we need for the rows writer.
type WriterOpts struct {
Compress bool
// FrameThresh is the minimum frame size in uncompressed bytes.
FrameThresh int
Rows bool
FrameCap uint64
Rows bool
}

const DefaultFrameCap = 20_000

// ColumnWriter implements the vio.Pusher interface. A Pusher creates a vector
// BSUP object from a stream of vector.Any.
type ColumnWriter struct {
writer io.WriteCloser
dynamic *vbuild.DynamicBuilder
sctx *super.Context
fuser fuser
size uint64
ctrl *RowWriter
writer io.WriteCloser
dynamic *vbuild.DynamicBuilder
sctx *super.Context
fuser fuser
size uint64
ctrl *RowWriter
framecap uint64
}

var _ vio.Pusher = (*ColumnWriter)(nil)

func NewColumnWriter(w io.WriteCloser) *ColumnWriter {
return NewColumnWriterWithCap(w, DefaultFrameCap)
}

func NewColumnWriterWithCap(w io.WriteCloser, frameCap uint64) *ColumnWriter {
Comment on lines 39 to +43

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: I think it's tidier to have just one constructor for which the zero value for a parameter gets you a reasonable default.

sctx := super.NewContext()
return &ColumnWriter{
writer: w,
dynamic: vbuild.NewDynamicBuilder(),
sctx: sctx,
fuser: newFuser(sctx),
writer: w,
dynamic: vbuild.NewDynamicBuilder(),
sctx: sctx,
fuser: newFuser(sctx),
framecap: frameCap,
}
}

func NewWriterWithOpts(w io.WriteCloser, opt WriterOpts) vio.PushCloser {
if opt.Rows {
return NewRowWriter(w)
writer := NewRowWriter(w)
writer.framecap = uint32(opt.FrameCap)
return writer
}
return NewColumnWriter(w)
writer := NewColumnWriter(w)
writer.framecap = opt.FrameCap
return writer
}

func (c *ColumnWriter) Close() error {
Expand Down Expand Up @@ -87,7 +94,7 @@ func (c *ColumnWriter) WriteSuperFrame(vec vector.Any) (uint64, error) {
func (c *ColumnWriter) Push(vec vector.Any) error {
if vec.Len() != 0 {
c.dynamic.Write(vec)
if c.dynamic.Len() >= maxFrameSize {
if c.dynamic.Len() >= uint32(c.framecap) {
return c.pushFrame()
}
}
Expand Down Expand Up @@ -202,6 +209,7 @@ type RowWriter struct {
size uint64
bytes []byte
len uint32
framecap uint32
}

var _ vio.Pusher = (*RowWriter)(nil)
Expand All @@ -215,6 +223,7 @@ func NewRowWriter(w io.WriteCloser) *RowWriter {
sctx: sctx,
fuser: newFuser(sctx),
superfuser: newFuser(sctx),
framecap: DefaultFrameCap,
}
}

Expand All @@ -237,7 +246,7 @@ func (r *RowWriter) Push(vec vector.Any) error {
}
}
r.len += vec.Len()
if r.len >= maxFrameSize {
if r.len >= r.framecap {
return r.pushFrame(false)
}
return nil
Expand All @@ -253,7 +262,7 @@ func (r *RowWriter) write(val super.Value, ctrl bool) error {
r.fuser.fuse(typ)
r.superfuser.fuse(typ)
r.len++
if r.len >= maxFrameSize || ctrl {
if r.len >= r.framecap || ctrl {
return r.pushFrame(ctrl)
}
return nil
Expand Down
9 changes: 9 additions & 0 deletions bsup/ztests/framecap.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
script: |
seq 100 | super -inputcap 5 -framecap 5 -o out.bsup -

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: Put a prefix on this flag since it only affects BSUP.

Suggested change
seq 100 | super -inputcap 5 -framecap 5 -o out.bsup -
seq 100 | super -inputcap 5 -bsup.framecap 5 -o out.bsup -

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Let's discuss. I'm gonna leave this for now and we can visit flags names overall.

super dev bsup out.bsup | super -s -c "count() by nameof(typeof(this)) | ? Header or Footer | sort count" -

outputs:
- name: stdout
data: |
{nameof:"SuperFooter",count:1}
{nameof:"ColumnHeader",count:20}
1 change: 1 addition & 0 deletions cli/inputflags/flags.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ func (f *Flags) SetFlags(fs *flag.FlagSet) {
})
fs.BoolVar(&f.Dynamic, "dynamic", false, "disable static type checking of inputs")
fs.StringVar(&opts.Format, "i", "auto", "format of input data [auto,arrows,bsup,csv,json,line,parquet,sup,tsv,zeek]")
fs.IntVar(&opts.InputCap, "inputcap", 0, "limit size of batched units of input")

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: Put this after "i" so these remain ordered by flag name.

fs.BoolVar(&f.Static, "static", false, "force static type checking of inputs")
}

Expand Down
1 change: 1 addition & 0 deletions cli/outputflags/flags.go
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ func (f *Flags) Options() anyio.WriterOpts {

func (f *Flags) setFlags(fs *flag.FlagSet) {
fs.BoolVar(&f.color, "color", true, "enable/disable color formatting for -S and db text output")
fs.Uint64Var(&f.BSUP.FrameCap, "framecap", 10000, "number of values per BSUP frame")
fs.BoolVar(&f.CSV.NoHeader, "noheader", false, "omit header for CSV and TSV output")
fs.IntVar(&f.pretty, "pretty", 2,
"tab size to pretty print JSON and Super JSON output (0 for newline-delimited output")
Expand Down
15 changes: 8 additions & 7 deletions cmd/super/db/create/command.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,12 +5,12 @@ import (
"flag"
"fmt"

"github.com/superdb/super/bsup"
"github.com/superdb/super/cli/poolflags"
"github.com/superdb/super/cmd/super/db"
"github.com/superdb/super/db/data"
"github.com/superdb/super/order"
"github.com/superdb/super/pkg/charm"
"github.com/superdb/super/pkg/units"
)

var spec = &charm.Spec{
Expand All @@ -25,9 +25,10 @@ See https://superdb.org/command/db.html#super-db-create

type Command struct {
*db.Command
sortKey string
thresh units.Bytes
use bool
sortKey string
frameCap uint64
objectCap uint64
use bool
}

func init() {
Expand All @@ -36,8 +37,8 @@ func init() {

func New(parent charm.Command, f *flag.FlagSet) (charm.Command, error) {
c := &Command{Command: parent.(*db.Command)}
c.thresh = data.DefaultThreshold
f.Var(&c.thresh, "S", "target size of pool data objects, as '10MB' or '4GiB', etc.")
f.Uint64Var(&c.frameCap, "framecap", bsup.DefaultFrameCap, "target number of values BSUP frames")
f.Uint64Var(&c.objectCap, "objectcap", data.DefaultObjectCap, "target number of values in pool data objects")
f.BoolVar(&c.use, "use", false, "set created pool as the current pool")
f.StringVar(&c.sortKey, "orderby", "ts:desc", "pool key with optional :asc or :desc suffix to organize data in pool (cannot be changed)")
return c, nil
Expand All @@ -61,7 +62,7 @@ func (c *Command) Run(args []string) error {
return err
}
poolName := args[0]
id, err := db.CreatePool(ctx, poolName, sortKey, int64(c.thresh))
id, err := db.CreatePool(ctx, poolName, sortKey, c.objectCap, c.frameCap)
if err != nil {
return err
}
Expand Down
8 changes: 4 additions & 4 deletions cmd/super/db/internal/dbmanage/scan.go
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,7 @@ func scan(ctx context.Context, it *objectIterator, pool *pools.Config, runCh cha
if err != nil {
return err
}
if run.overlaps(o.Min, o.Max) || run.size+o.Size < pool.Threshold {
if run.overlaps(o.Min, o.Max) || run.len+o.Count < pool.ObjectCap {
run.add(o)
continue
}
Expand Down Expand Up @@ -119,7 +119,7 @@ type runBuilder struct {
span extent.Span
cmp expr.CompareFn
objects []*object
size int64
len uint64
}

func newRunBuilder() *runBuilder {
Expand All @@ -135,7 +135,7 @@ func (r *runBuilder) overlaps(first, last super.Value) bool {

func (r *runBuilder) add(o *object) {
r.objects = append(r.objects, o)
r.size += o.Size
r.len += o.Count
if r.span == nil {
r.span = extent.NewGeneric(o.Min, o.Max, r.cmp)
return
Expand All @@ -155,5 +155,5 @@ func (r *runBuilder) objectIDs() []uuid.UUID {
func (r *runBuilder) reset() {
r.span = nil
r.objects = r.objects[:0]
r.size = 0
r.len = 0
}
2 changes: 1 addition & 1 deletion cmd/super/ztests/terminal-output-format.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -15,4 +15,4 @@ outputs:
- name: stdout
data: "{a:1}\r\n\
{b:2}\r\n\
p1 xxx key ts order desc\r\n"
p1 xxx key ts order desc objectcap 1048576 framecap 20000\r\n"
2 changes: 1 addition & 1 deletion db/api/api.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@ type Interface interface {
Query(ctx context.Context, query []srcfiles.Input) (vio.Scanner, error)
PoolID(ctx context.Context, poolName string) (uuid.UUID, error)
CommitObject(ctx context.Context, poolID uuid.UUID, branchName string) (uuid.UUID, error)
CreatePool(context.Context, string, order.SortKeys, int64) (uuid.UUID, error)
CreatePool(context.Context, string, order.SortKeys, uint64, uint64) (uuid.UUID, error)
RemovePool(context.Context, uuid.UUID) error
RenamePool(context.Context, uuid.UUID, string) error
CreateBranch(ctx context.Context, pool uuid.UUID, name string, parent uuid.UUID) error
Expand Down
4 changes: 2 additions & 2 deletions db/api/local.go
Original file line number Diff line number Diff line change
Expand Up @@ -61,11 +61,11 @@ func (l *local) Root() *db.Root {
return l.db
}

func (l *local) CreatePool(ctx context.Context, name string, sortKeys order.SortKeys, thresh int64) (uuid.UUID, error) {
func (l *local) CreatePool(ctx context.Context, name string, sortKeys order.SortKeys, objectCap, frameCap uint64) (uuid.UUID, error) {
if name == "" {
return uuid.Nil(), errors.New("no pool name provided")
}
pool, err := l.db.CreatePool(ctx, name, sortKeys, thresh)
pool, err := l.db.CreatePool(ctx, name, sortKeys, objectCap, frameCap)
if err != nil {
return uuid.Nil(), err
}
Expand Down
5 changes: 3 additions & 2 deletions db/api/remote.go
Original file line number Diff line number Diff line change
Expand Up @@ -53,14 +53,15 @@ func (r *remote) CommitObject(ctx context.Context, poolID uuid.UUID, branchName
return res.Commit, err
}

func (r *remote) CreatePool(ctx context.Context, name string, sortKeys order.SortKeys, thresh int64) (uuid.UUID, error) {
func (r *remote) CreatePool(ctx context.Context, name string, sortKeys order.SortKeys, objectCap, frameCap uint64) (uuid.UUID, error) {
res, err := r.conn.CreatePool(ctx, api.PoolPostRequest{
Name: name,
SortKeys: api.SortKeys{
Order: sortKeys.Primary().Order,
Keys: field.List{sortKeys.Primary().Path},
},
Thresh: thresh,
ObjectCap: objectCap,
FrameCap: frameCap,
})
if err != nil {
return uuid.Nil(), err
Expand Down
2 changes: 1 addition & 1 deletion db/data/object.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ import (
)

const (
DefaultThreshold = 500 * 1024 * 1024
DefaultObjectCap = 1024 * 1024
)
Comment on lines 17 to 19

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit:

const DefaultObjectCap = 024 * 1024


// A FileKind is the first part of a file name, used to differentiate files
Expand Down
16 changes: 11 additions & 5 deletions db/pools/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package pools
import (
"uuid"

"github.com/superdb/super/bsup"
"github.com/superdb/super/db/data"
"github.com/superdb/super/db/journal"
"github.com/superdb/super/order"
Expand All @@ -16,24 +17,29 @@ type Config struct {
Name string `super:"name"`
ID uuid.UUID `super:"id"`
SortKeys order.SortKeys `super:"layout"`
Threshold int64 `super:"threshold"`
ObjectCap uint64 `super:"objectcap"`
FrameCap uint64 `super:"framecap"`
}

var _ journal.Entry = (*Config)(nil)

func NewConfig(name string, sortKeys order.SortKeys, thresh int64) *Config {
func NewConfig(name string, sortKeys order.SortKeys, objectCap, frameCap uint64) *Config {
if sortKeys.IsNil() {
sortKeys = order.SortKeys{order.NewSortKey(order.Desc, field.Dotted("ts"))}
}
if thresh == 0 {
thresh = data.DefaultThreshold
if objectCap == 0 {
objectCap = data.DefaultObjectCap
}
if frameCap == 0 {
frameCap = bsup.DefaultFrameCap
}
return &Config{
Ts: nano.Now(),
Name: name,
ID: uuid.NewV7(),
SortKeys: sortKeys,
Threshold: thresh,
ObjectCap: objectCap,
FrameCap: frameCap,
}
}

Expand Down
14 changes: 9 additions & 5 deletions db/root.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (

arc "github.com/hashicorp/golang-lru/arc/v2"
"github.com/superdb/super"
"github.com/superdb/super/bsup"
"github.com/superdb/super/bsupbytes"
"github.com/superdb/super/compiler/dag"
"github.com/superdb/super/db/branches"
Expand Down Expand Up @@ -315,20 +316,23 @@ func (r *Root) RenamePool(ctx context.Context, id uuid.UUID, newName string) err
return r.pools.Rename(ctx, id, newName)
}

func (r *Root) CreatePool(ctx context.Context, name string, sortKeys order.SortKeys, thresh int64) (*Pool, error) {
func (r *Root) CreatePool(ctx context.Context, name string, sortKeys order.SortKeys, objectCap, frameCap uint64) (*Pool, error) {
if name == "HEAD" {
return nil, fmt.Errorf("pool cannot be named %q", name)
}
if r.pools.LookupByName(ctx, name) != nil {
return nil, fmt.Errorf("%s: %w", name, pools.ErrExists)
}
if thresh == 0 {
thresh = data.DefaultThreshold
}
if len(sortKeys) > 1 {
return nil, errors.New("multiple pool keys not supported")
}
config := pools.NewConfig(name, sortKeys, thresh)
if objectCap == 0 {
objectCap = data.DefaultObjectCap
}
if frameCap == 0 {
frameCap = bsup.DefaultFrameCap
}
config := pools.NewConfig(name, sortKeys, objectCap, frameCap)
if err := CreatePool(ctx, r.engine, r.logger, r.path, config); err != nil {
return nil, err
}
Expand Down
Loading
Loading