mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-08-28 02:44:59 +00:00
[FB-1822] Change dataframe disk format to Arrow from parquet (#2376)
* change default backend to arrow file format instead of parquet
(cherry picked from commit 6d96a1474e)
This commit is contained in:
parent
2993ab5ca1
commit
24372a9405
9 changed files with 243 additions and 74 deletions
84
apply.go
84
apply.go
|
|
@ -13,8 +13,6 @@ import (
|
|||
"github.com/apache/arrow/go/v10/arrow"
|
||||
"github.com/apache/arrow/go/v10/arrow/array"
|
||||
"github.com/apache/arrow/go/v10/arrow/memory"
|
||||
"github.com/apache/arrow/go/v10/parquet/file"
|
||||
"github.com/apache/arrow/go/v10/parquet/pqarrow"
|
||||
"github.com/gomem/gomem/pkg/dataframe"
|
||||
"github.com/featurebasedb/featurebase/v3/pql"
|
||||
"github.com/featurebasedb/featurebase/v3/tracing"
|
||||
|
|
@ -220,12 +218,14 @@ func (e *executor) executeApplyShard(ctx context.Context, qcx *Qcx, index string
|
|||
if idx == nil {
|
||||
return nil, newNotFoundError(ErrIndexNotFound, index)
|
||||
}
|
||||
|
||||
fname := idx.GetDataFramePath(shard)
|
||||
if _, err := os.Stat(fname + ".parquet"); os.IsNotExist(err) {
|
||||
|
||||
if !e.dataFrameExists(fname) {
|
||||
return value.NewVector([]value.Value{}), nil
|
||||
}
|
||||
|
||||
table, err := readTableParquet(fname)
|
||||
table, err := e.getDataTable(ctx, fname, pool)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -254,57 +254,20 @@ func (e *executor) executeApplyShard(ctx context.Context, qcx *Qcx, index string
|
|||
return context.Global("_"), nil
|
||||
}
|
||||
|
||||
func readTableParquet(filename string) (arrow.Table, error) {
|
||||
r, err := os.Open(filename + ".parquet")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
pf, err := file.NewParquetReader(r)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
reader, err := pqarrow.NewFileReader(pf, pqarrow.ArrowReadProperties{}, memory.DefaultAllocator)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return reader.ReadTable(context.Background())
|
||||
}
|
||||
|
||||
func readTableParquetCtx(ctx context.Context, filename string, mem memory.Allocator) (arrow.Table, error) {
|
||||
r, err := os.Open(filename + ".parquet")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
pf, err := file.NewParquetReader(r)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
reader, err := pqarrow.NewFileReader(pf, pqarrow.ArrowReadProperties{}, mem)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return reader.ReadTable(ctx)
|
||||
}
|
||||
|
||||
// ///////////////////////////////////////////////////////
|
||||
// all the ingest supporting functions
|
||||
// ///////////////////////////////////////////////////////
|
||||
|
||||
func NewShardFile(name string) (*ShardFile, error) {
|
||||
if _, err := os.Stat(name + ".parquet"); os.IsNotExist(err) {
|
||||
return &ShardFile{dest: name}, nil
|
||||
func NewShardFile(ctx context.Context, name string, mem memory.Allocator, e *executor) (*ShardFile, error) {
|
||||
if !e.dataFrameExists(name) {
|
||||
return &ShardFile{dest: name, executor: e}, nil
|
||||
}
|
||||
|
||||
// else read in existing
|
||||
table, err := readTableParquet(name)
|
||||
table, err := e.getDataTable(ctx, name, mem)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &ShardFile{table: table, schema: table.Schema(), dest: name}, nil
|
||||
return &ShardFile{table: table, schema: table.Schema(), dest: name, executor: e}, nil
|
||||
}
|
||||
|
||||
type NameType struct {
|
||||
|
|
@ -349,6 +312,7 @@ type ShardFile struct {
|
|||
added int64
|
||||
columns []interface{}
|
||||
dest string
|
||||
executor *executor
|
||||
}
|
||||
|
||||
func compareSchema(s1, s2 *arrow.Schema) bool {
|
||||
|
|
@ -427,7 +391,7 @@ func (sf *ShardFile) Process(cs *ChangesetRequest) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return os.Rename(rtemp+".parquet", sf.dest+".parquet")
|
||||
return os.Rename(rtemp+sf.executor.TableExtension(), sf.dest+sf.executor.TableExtension())
|
||||
}
|
||||
|
||||
func (sf *ShardFile) process(cs *ChangesetRequest) error {
|
||||
|
|
@ -521,22 +485,9 @@ func (sf *ShardFile) Save(name string) error {
|
|||
}
|
||||
}
|
||||
rec := array.NewRecord(sf.schema, parts, sf.beforeRows+sf.added)
|
||||
df, err := dataframe.NewDataFrameFromRecord(mem, rec)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// confirm change
|
||||
w, err := os.Create(name + ".parquet")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
table := array.NewTableFromRecords(sf.schema, []arrow.Record{rec})
|
||||
|
||||
err = df.ToParquet(w, 1024)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
w.Close()
|
||||
return nil
|
||||
return sf.executor.SaveTable(name, table, mem)
|
||||
}
|
||||
|
||||
// TODO(twg) 2022/10/03 Not a huge fan of the global variable will look at adding to executor structure
|
||||
|
|
@ -573,7 +524,8 @@ func (api *API) ApplyDataframeChangeset(ctx context.Context, index string, cs *C
|
|||
mu := getDataframeWritelock(shard)
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
shardFile, err := NewShardFile(fname)
|
||||
mem := memory.NewGoAllocator()
|
||||
shardFile, err := NewShardFile(ctx, fname, mem, api.server.executor)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -600,14 +552,16 @@ func (api *API) GetDataframeSchema(ctx context.Context, indexName string) (inter
|
|||
dir, _ := os.Open(base)
|
||||
files, _ := dir.Readdir(0)
|
||||
parts := make([]column, 0)
|
||||
mem := memory.NewGoAllocator()
|
||||
for i := range files {
|
||||
file := files[i]
|
||||
name := file.Name()
|
||||
if strings.HasSuffix(name, ".parquet") {
|
||||
if api.server.executor.IsDataframeFile(name) {
|
||||
// strip off the parquet extenison
|
||||
name = strings.TrimSuffix(name, filepath.Ext(name))
|
||||
// read the parquet file and extract the schema
|
||||
table, err := readTableParquet(filepath.Join(base, name))
|
||||
fname := filepath.Join(base, name)
|
||||
table, err := api.server.executor.getDataTable(ctx, fname, mem)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
|
|||
151
arrow.go
151
arrow.go
|
|
@ -5,12 +5,18 @@ import (
|
|||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
"github.com/apache/arrow/go/v10/arrow"
|
||||
"github.com/apache/arrow/go/v10/arrow/array"
|
||||
"github.com/apache/arrow/go/v10/arrow/ipc"
|
||||
"github.com/apache/arrow/go/v10/arrow/memory"
|
||||
"github.com/apache/arrow/go/v10/parquet"
|
||||
"github.com/apache/arrow/go/v10/parquet/file"
|
||||
"github.com/apache/arrow/go/v10/parquet/pqarrow"
|
||||
"github.com/gomem/gomem/pkg/dataframe"
|
||||
"github.com/featurebasedb/featurebase/v3/pql"
|
||||
"github.com/featurebasedb/featurebase/v3/tracing"
|
||||
|
|
@ -371,12 +377,14 @@ func (e *executor) executeArrowShard(ctx context.Context, qcx *Qcx, index string
|
|||
if idx == nil {
|
||||
return nil, newNotFoundError(ErrIndexNotFound, index)
|
||||
}
|
||||
|
||||
fname := idx.GetDataFramePath(shard)
|
||||
if _, err := os.Stat(fname + ".parquet"); os.IsNotExist(err) {
|
||||
|
||||
if !e.dataFrameExists(fname) {
|
||||
return &basicTable{name: name}, nil
|
||||
}
|
||||
|
||||
table, err := readTableParquetCtx(context.TODO(), fname, pool)
|
||||
table, err := e.getDataTable(ctx, fname, pool)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "arrow readTableParquet")
|
||||
}
|
||||
|
|
@ -403,3 +411,142 @@ func (e *executor) executeArrowShard(ctx context.Context, qcx *Qcx, index string
|
|||
table.Retain()
|
||||
return &basicTable{resolver: resolver, table: table, filtered: filter != nil, name: name}, nil
|
||||
}
|
||||
|
||||
func (e *executor) dataFrameExists(fname string) bool {
|
||||
if e.typeIsParquet() {
|
||||
if _, err := os.Stat(fname + ".parquet"); os.IsNotExist(err) {
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
if _, err := os.Stat(fname + ".arrow"); os.IsNotExist(err) {
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func (e *executor) getDataTable(ctx context.Context, fname string, mem memory.Allocator) (arrow.Table, error) {
|
||||
if e.typeIsParquet() {
|
||||
table, err := readTableParquetCtx(ctx, fname, mem)
|
||||
return table, err
|
||||
}
|
||||
return readTableArrow(fname, mem)
|
||||
}
|
||||
|
||||
func (e *executor) typeIsParquet() bool {
|
||||
return e.datafameUseParquet
|
||||
}
|
||||
|
||||
func (e *executor) IsDataframeFile(name string) bool {
|
||||
if e.typeIsParquet() {
|
||||
return strings.HasSuffix(name, ".parquet")
|
||||
}
|
||||
return strings.HasSuffix(name, ".arrow")
|
||||
}
|
||||
|
||||
func (e *executor) SaveTable(name string, table arrow.Table, mem memory.Allocator) error {
|
||||
if e.typeIsParquet() {
|
||||
return writeTableParquet(table, name)
|
||||
}
|
||||
return writeTableArrow(table, name, mem)
|
||||
}
|
||||
|
||||
func (e *executor) TableExtension() string {
|
||||
if e.typeIsParquet() {
|
||||
return ".parquet"
|
||||
}
|
||||
return ".arrow"
|
||||
}
|
||||
|
||||
func readTableArrow(filename string, mem memory.Allocator) (arrow.Table, error) {
|
||||
r, err := os.Open(filename + ".arrow")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
rr, err := ipc.NewFileReader(r, ipc.WithAllocator(mem))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rr.Close()
|
||||
records := make([]arrow.Record, rr.NumRecords(), rr.NumRecords())
|
||||
i := 0
|
||||
for {
|
||||
rec, err := rr.Read()
|
||||
if err == io.EOF {
|
||||
break
|
||||
} else if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
records[i] = rec
|
||||
i++
|
||||
}
|
||||
records = records[:i]
|
||||
table := array.NewTableFromRecords(rr.Schema(), records)
|
||||
return table, nil
|
||||
}
|
||||
|
||||
func readTableParquetCtx(ctx context.Context, filename string, mem memory.Allocator) (arrow.Table, error) {
|
||||
r, err := os.Open(filename + ".parquet")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer r.Close()
|
||||
|
||||
pf, err := file.NewParquetReader(r)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
reader, err := pqarrow.NewFileReader(pf, pqarrow.ArrowReadProperties{}, mem)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return reader.ReadTable(ctx)
|
||||
}
|
||||
|
||||
func writeTableParquet(table arrow.Table, filename string) error {
|
||||
f, err := os.Create(filename + ".parquet")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer f.Close()
|
||||
props := parquet.NewWriterProperties(parquet.WithDictionaryDefault(false))
|
||||
arrProps := pqarrow.DefaultWriterProps()
|
||||
chunkSize := 10 * 1024 * 1024
|
||||
err = pqarrow.WriteTable(table, f, int64(chunkSize), props, arrProps)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
f.Sync()
|
||||
return nil
|
||||
}
|
||||
|
||||
func writeTableArrow(table arrow.Table, filename string, mem memory.Allocator) error {
|
||||
f, err := os.Create(filename + ".arrow")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer f.Close()
|
||||
writer, err := ipc.NewFileWriter(f, ipc.WithAllocator(mem), ipc.WithSchema(table.Schema()))
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
chunkSize := int64(0)
|
||||
tr := array.NewTableReader(table, chunkSize)
|
||||
defer tr.Release()
|
||||
n := 0
|
||||
for tr.Next() {
|
||||
arec := tr.Record()
|
||||
err = writer.Write(arec)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
n++
|
||||
}
|
||||
err = writer.Close()
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
f.Sync()
|
||||
return nil
|
||||
}
|
||||
|
|
|
|||
53
arrow_test.go
Normal file
53
arrow_test.go
Normal file
|
|
@ -0,0 +1,53 @@
|
|||
// Copyright 2021 Molecula Corp. All rights reserved.
|
||||
package pilosa
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/hex"
|
||||
"math/rand"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
|
||||
"github.com/apache/arrow/go/v10/arrow"
|
||||
"github.com/apache/arrow/go/v10/arrow/array"
|
||||
"github.com/apache/arrow/go/v10/arrow/memory"
|
||||
)
|
||||
|
||||
func TempFileName(prefix string) string {
|
||||
randBytes := make([]byte, 16)
|
||||
rand.Read(randBytes)
|
||||
return filepath.Join(os.TempDir(), prefix+hex.EncodeToString(randBytes))
|
||||
}
|
||||
|
||||
func Test_TableParquet(t *testing.T) {
|
||||
// create a arrow table
|
||||
schema := arrow.NewSchema(
|
||||
[]arrow.Field{
|
||||
{Name: "num", Type: arrow.PrimitiveTypes.Float64},
|
||||
},
|
||||
nil, // no metadata
|
||||
)
|
||||
mem := memory.NewGoAllocator()
|
||||
b := array.NewRecordBuilder(mem, schema)
|
||||
defer b.Release()
|
||||
b.Field(0).(*array.Float64Builder).AppendValues([]float64{1.0, 1.5, 2.0}, nil)
|
||||
table := array.NewTableFromRecords(schema, []arrow.Record{b.NewRecord()})
|
||||
defer table.Release()
|
||||
fileName := TempFileName("pq-")
|
||||
// save it as a parquet file
|
||||
err := writeTableParquet(table, fileName)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer os.Remove(fileName)
|
||||
|
||||
// read it back in and compare the result
|
||||
got, err := readTableParquetCtx(context.Background(), fileName, mem)
|
||||
if err != nil {
|
||||
t.Fatalf("readTableParquetCtx() error = %v", err)
|
||||
}
|
||||
if got.NumCols() != table.NumCols() {
|
||||
t.Errorf("got:%v expected:%v", got.NumCols(), table.NumCols())
|
||||
}
|
||||
}
|
||||
|
|
@ -144,6 +144,7 @@ func serverFlagSet(srv *server.Config, prefix string) *pflag.FlagSet {
|
|||
flags.BoolVar(&srv.DataDog.BlockProfile, pre("datadog.block-profile"), false, "golang pprof goroutine ")
|
||||
|
||||
flags.BoolVar(&srv.Dataframe.Enable, pre("dataframe.enable"), false, "EXPERIMENTAL enable support for Apply and Arrow")
|
||||
flags.BoolVar(&srv.Dataframe.UseParquet, pre("dataframe.use-parquet"), false, "EXPERIMENTAL use parquet for file format")
|
||||
|
||||
return flags
|
||||
}
|
||||
|
|
|
|||
|
|
@ -77,7 +77,8 @@ type executor struct {
|
|||
maxMemory int64
|
||||
|
||||
// Temporary flag to be removed when stablized
|
||||
dataframeEnabled bool
|
||||
dataframeEnabled bool
|
||||
datafameUseParquet bool
|
||||
}
|
||||
|
||||
// executorOption is a functional option type for pilosa.executor
|
||||
|
|
|
|||
|
|
@ -3860,7 +3860,7 @@ func (h *Handler) handlePostDataframeRestore(w http.ResponseWriter, r *http.Requ
|
|||
http.Error(w, fmt.Sprintf("Index %s Not Found", indexName), http.StatusNotFound)
|
||||
return
|
||||
}
|
||||
filename := idx.GetDataFramePath(shard) + ".parquet"
|
||||
filename := idx.GetDataFramePath(shard) + h.api.server.executor.TableExtension()
|
||||
dest, err := os.Create(filename)
|
||||
if err != nil {
|
||||
http.Error(w, fmt.Sprintf("failed to create restore dataframe shard %v %v err:%v", indexName, shard, err), http.StatusBadRequest)
|
||||
|
|
@ -4184,7 +4184,7 @@ func (h *Handler) handleGetDataframe(w http.ResponseWriter, r *http.Request) {
|
|||
http.Error(w, fmt.Sprintf("Index %s Not Found", indexName), http.StatusNotFound)
|
||||
return
|
||||
}
|
||||
filename := idx.GetDataFramePath(shard) + ".parquet"
|
||||
filename := idx.GetDataFramePath(shard) + h.api.server.executor.TableExtension()
|
||||
http.ServeFile(w, r, filename)
|
||||
}
|
||||
|
||||
|
|
|
|||
11
server.go
11
server.go
|
|
@ -99,7 +99,8 @@ type Server struct { // nolint: maligned
|
|||
|
||||
serverlessStorage *daxstorage.ResourceManager
|
||||
|
||||
dataframeEnabled bool
|
||||
dataframeEnabled bool
|
||||
dataframeUseParquet bool
|
||||
}
|
||||
|
||||
type ExecutionPlannerFn func(executor Executor, api *API, sql string) sql3.CompilePlanner
|
||||
|
|
@ -457,6 +458,13 @@ func OptServerIsDataframeEnabled(is bool) ServerOption {
|
|||
}
|
||||
}
|
||||
|
||||
func OptServerDataframeUseParquet(is bool) ServerOption {
|
||||
return func(s *Server) error {
|
||||
s.dataframeUseParquet = is
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
// NewServer returns a new instance of Server.
|
||||
func NewServer(opts ...ServerOption) (*Server, error) {
|
||||
cluster := newCluster()
|
||||
|
|
@ -528,6 +536,7 @@ func NewServer(opts ...ServerOption) (*Server, error) {
|
|||
}
|
||||
s.executor = newExecutor(executorOpts...)
|
||||
s.executor.dataframeEnabled = s.dataframeEnabled
|
||||
s.executor.datafameUseParquet = s.dataframeUseParquet
|
||||
|
||||
path, err := expandDirName(s.dataDir)
|
||||
if err != nil {
|
||||
|
|
|
|||
|
|
@ -252,7 +252,8 @@ type Config struct {
|
|||
Auth Auth
|
||||
|
||||
Dataframe struct {
|
||||
Enable bool `toml:"enable"`
|
||||
Enable bool `toml:"enable"`
|
||||
UseParquet bool `toml:"use-parquet"`
|
||||
} `toml:"dataframe"`
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -195,8 +195,10 @@ const (
|
|||
|
||||
// we want to set resource limits *exactly once*, and then be able
|
||||
// to report on whether or not that succeeded.
|
||||
var setupResourceLimitsOnce sync.Once
|
||||
var setupResourceLimitsErr error
|
||||
var (
|
||||
setupResourceLimitsOnce sync.Once
|
||||
setupResourceLimitsErr error
|
||||
)
|
||||
|
||||
// doSetupResourceLimits is the function which actually does the
|
||||
// resource limit setup, possibly yielding an error. it's a Command
|
||||
|
|
@ -591,6 +593,7 @@ func (m *Command) setupServer() error {
|
|||
pilosa.OptServerExecutionPlannerFn(executionPlannerFn),
|
||||
pilosa.OptServerServerlessStorage(m.serverlessStorage),
|
||||
pilosa.OptServerIsDataframeEnabled(m.Config.Dataframe.Enable),
|
||||
pilosa.OptServerDataframeUseParquet(m.Config.Dataframe.UseParquet),
|
||||
}
|
||||
|
||||
if m.isComputeNode {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue