Improve CSV & JSON decoders for larger datasets
The built in buffer had a strange interaction with the CSV library. Changed both decoders for consistency.
This commit is contained in:
+24
-42
@@ -1,7 +1,6 @@
|
||||
package csv
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/csv"
|
||||
"io"
|
||||
@@ -18,11 +17,11 @@ type (
|
||||
|
||||
// csv decoder wrapper for additional bits
|
||||
reader struct {
|
||||
c *csv.Reader
|
||||
header []string
|
||||
count uint64
|
||||
|
||||
buff bytes.Buffer
|
||||
cache []map[string]string
|
||||
cacheIndex int
|
||||
}
|
||||
)
|
||||
|
||||
@@ -52,19 +51,12 @@ func (y *decoder) CanDecodeExt(ext string) bool {
|
||||
|
||||
// Decode decodes the given io.Reader into a generic resource dataset
|
||||
func (c *decoder) Decode(ctx context.Context, r io.Reader, do *envoy.DecoderOpts) ([]resource.Interface, error) {
|
||||
cr := &reader{}
|
||||
var buff bytes.Buffer
|
||||
|
||||
// So we can reset to the start of the reader
|
||||
err := cr.prepare(io.TeeReader(r, &buff))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
cr := &reader{
|
||||
cacheIndex: 0,
|
||||
cache: make([]map[string]string, 0, 1000),
|
||||
}
|
||||
|
||||
cr.c = csv.NewReader(io.TeeReader(&buff, &cr.buff))
|
||||
|
||||
// The first one is a header, so let's just get rid of it
|
||||
_, err = cr.c.Read()
|
||||
err := cr.prepare(r)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -72,6 +64,10 @@ func (c *decoder) Decode(ctx context.Context, r io.Reader, do *envoy.DecoderOpts
|
||||
return []resource.Interface{resource.NewResourceDataset(do.Name, cr)}, nil
|
||||
}
|
||||
|
||||
// The prepare step caches all of the rows into a cache array that will be used
|
||||
// to access already visited CSV rows.
|
||||
//
|
||||
// @todo implement FS cache file for larger files
|
||||
func (cr *reader) prepare(r io.Reader) (err error) {
|
||||
cReader := csv.NewReader(r)
|
||||
|
||||
@@ -81,19 +77,24 @@ func (cr *reader) prepare(r io.Reader) (err error) {
|
||||
return err
|
||||
}
|
||||
|
||||
// Entry count
|
||||
for {
|
||||
_, err := cReader.Read()
|
||||
aux := make(map[string]string)
|
||||
rr, err := cReader.Read()
|
||||
if err == io.EOF {
|
||||
break
|
||||
return nil
|
||||
} else if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Entry count
|
||||
cr.count++
|
||||
}
|
||||
|
||||
return nil
|
||||
// Cache entry
|
||||
for i, h := range cr.header {
|
||||
aux[h] = rr[i]
|
||||
}
|
||||
cr.cache = append(cr.cache, aux)
|
||||
}
|
||||
}
|
||||
|
||||
// Fields returns every available field in this dataset
|
||||
@@ -102,37 +103,18 @@ func (cr *reader) Fields() []string {
|
||||
}
|
||||
|
||||
func (cr *reader) Reset() error {
|
||||
if len(cr.buff.Bytes()) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
buff := cr.buff
|
||||
cr.buff = bytes.Buffer{}
|
||||
|
||||
cr.c = csv.NewReader(io.TeeReader(&buff, &cr.buff))
|
||||
// The first one is a header, so let's just get rid of it
|
||||
_, err := cr.c.Read()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
cr.cacheIndex = 0
|
||||
return nil
|
||||
}
|
||||
|
||||
// Next returns the field: value mapping for the next row
|
||||
func (cr *reader) Next() (map[string]string, error) {
|
||||
mr := make(map[string]string)
|
||||
rr, err := cr.c.Read()
|
||||
if err == io.EOF {
|
||||
if cr.cacheIndex >= len(cr.cache) {
|
||||
return nil, nil
|
||||
} else if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
for i, h := range cr.header {
|
||||
mr[h] = rr[i]
|
||||
}
|
||||
|
||||
mr := cr.cache[cr.cacheIndex]
|
||||
cr.cacheIndex++
|
||||
return mr, nil
|
||||
}
|
||||
|
||||
|
||||
+20
-29
@@ -18,11 +18,11 @@ type (
|
||||
|
||||
// json decoder wrapper for additional bits
|
||||
reader struct {
|
||||
d *json.Decoder
|
||||
header []string
|
||||
count uint64
|
||||
|
||||
buff bytes.Buffer
|
||||
cache []map[string]string
|
||||
cacheIndex int
|
||||
}
|
||||
)
|
||||
|
||||
@@ -56,20 +56,23 @@ func (d *decoder) CanDecodeExt(ext string) bool {
|
||||
|
||||
// Decode decodes the given io.Reader into a generic resource dataset
|
||||
func (d *decoder) Decode(ctx context.Context, r io.Reader, do *envoy.DecoderOpts) ([]resource.Interface, error) {
|
||||
jr := &reader{}
|
||||
var buff bytes.Buffer
|
||||
jr := &reader{
|
||||
cacheIndex: 0,
|
||||
cache: make([]map[string]string, 0, 1000),
|
||||
}
|
||||
|
||||
// So we can reset to the start of the reader
|
||||
err := jr.prepare(io.TeeReader(r, &buff))
|
||||
err := jr.prepare(r)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
jr.d = json.NewDecoder(io.TeeReader(&buff, &jr.buff))
|
||||
|
||||
return []resource.Interface{resource.NewResourceDataset(do.Name, jr)}, nil
|
||||
}
|
||||
|
||||
// The prepare step caches all of the rows into a cache array that will be used
|
||||
// to access already visited json rows.
|
||||
//
|
||||
// @todo implement FS cache file for larger files
|
||||
func (jr *reader) prepare(r io.Reader) (err error) {
|
||||
jReader := json.NewDecoder(r)
|
||||
|
||||
@@ -79,8 +82,8 @@ func (jr *reader) prepare(r io.Reader) (err error) {
|
||||
hx := make(map[string]bool)
|
||||
jr.header = make([]string, 0, 100)
|
||||
|
||||
aux := make(map[string]interface{})
|
||||
for jReader.More() {
|
||||
aux := make(map[string]string)
|
||||
err = jReader.Decode(&aux)
|
||||
if err == io.EOF {
|
||||
break
|
||||
@@ -89,7 +92,6 @@ func (jr *reader) prepare(r io.Reader) (err error) {
|
||||
}
|
||||
|
||||
// Get all the header fields
|
||||
// @todo do we want to preserve order?
|
||||
for h := range aux {
|
||||
if !hx[h] {
|
||||
jr.header = append(jr.header, h)
|
||||
@@ -97,22 +99,18 @@ func (jr *reader) prepare(r io.Reader) (err error) {
|
||||
}
|
||||
}
|
||||
|
||||
// Entry count
|
||||
jr.count++
|
||||
|
||||
// Cache entry
|
||||
jr.cache = append(jr.cache, aux)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (jr *reader) Reset() error {
|
||||
if len(jr.buff.Bytes()) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
buff := jr.buff
|
||||
jr.buff = bytes.Buffer{}
|
||||
|
||||
jr.d = json.NewDecoder(io.TeeReader(&buff, &jr.buff))
|
||||
|
||||
jr.cacheIndex = 0
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -123,19 +121,12 @@ func (jr *reader) Fields() []string {
|
||||
|
||||
// Next returns the field: value mapping for the next row
|
||||
func (jr *reader) Next() (map[string]string, error) {
|
||||
// It's over
|
||||
if !jr.d.More() {
|
||||
if jr.cacheIndex >= len(jr.cache) {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
mr := make(map[string]string)
|
||||
err := jr.d.Decode(&mr)
|
||||
if err == io.EOF {
|
||||
return nil, nil
|
||||
} else if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
mr := jr.cache[jr.cacheIndex]
|
||||
jr.cacheIndex++
|
||||
return mr, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -108,7 +108,7 @@ func NewStoreEncoder(s store.Storer, cfg *EncoderConfig) envoy.PrepareEncoder {
|
||||
// It initializes and prepares the resource state for each provided resource
|
||||
func (se *storeEncoder) Prepare(ctx context.Context, ee ...*envoy.ResourceState) (err error) {
|
||||
f := func(rs resourceState, ers *envoy.ResourceState) error {
|
||||
err = rs.Prepare(ctx, se.makePayload(ctx, ers))
|
||||
err = rs.Prepare(ctx, se.makePayload(ctx, se.s, ers))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -178,7 +178,7 @@ func (se *storeEncoder) Encode(ctx context.Context, p envoy.Provider) error {
|
||||
if state == nil {
|
||||
err = ErrResourceStateUndefined
|
||||
} else {
|
||||
err = state.Encode(ctx, se.makePayload(ctx, ers))
|
||||
err = state.Encode(ctx, se.makePayload(ctx, s, ers))
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
@@ -188,9 +188,9 @@ func (se *storeEncoder) Encode(ctx context.Context, p envoy.Provider) error {
|
||||
})
|
||||
}
|
||||
|
||||
func (se *storeEncoder) makePayload(ctx context.Context, ers *envoy.ResourceState) *payload {
|
||||
func (se *storeEncoder) makePayload(ctx context.Context, s store.Storer, ers *envoy.ResourceState) *payload {
|
||||
return &payload{
|
||||
s: se.s,
|
||||
s: s,
|
||||
state: ers,
|
||||
composeAccessControl: service.AccessControl(rbac.Global()),
|
||||
invokerID: auth.GetIdentityFromContext(ctx).Identity(),
|
||||
|
||||
@@ -86,6 +86,72 @@ func TestDataShaping(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestDataShaping_large(t *testing.T) {
|
||||
var (
|
||||
ctx = auth.SetSuperUserContext(context.Background())
|
||||
s = initStore(ctx, t)
|
||||
err error
|
||||
|
||||
cases = []string{
|
||||
"csv_large",
|
||||
"jsonl_large",
|
||||
}
|
||||
)
|
||||
|
||||
ni := uint64(10)
|
||||
su.NextID = func() uint64 {
|
||||
ni++
|
||||
return ni
|
||||
}
|
||||
|
||||
for _, c := range cases {
|
||||
t.Run(fmt.Sprintf("record shaping; data_shaping; large datasets/%s", c), func(t *testing.T) {
|
||||
var (
|
||||
req = require.New(t)
|
||||
)
|
||||
|
||||
truncateStore(ctx, s, t)
|
||||
err = collect(
|
||||
err,
|
||||
storeRole(ctx, s, 1, "everyone"),
|
||||
storeRole(ctx, s, 2, "admins"),
|
||||
)
|
||||
if err != nil {
|
||||
t.Fatal(err.Error())
|
||||
}
|
||||
|
||||
nn, err := decodeDirectory(ctx, path.Join("data_shaping", c))
|
||||
req.NoError(err)
|
||||
|
||||
crs := resource.ComposeRecordShaper()
|
||||
nn, err = resource.Shape(nn, crs)
|
||||
req.NoError(err)
|
||||
|
||||
req.NoError(encode(ctx, s, nn))
|
||||
|
||||
ns, err := store.LookupComposeNamespaceBySlug(ctx, s, "ns1")
|
||||
req.NotNil(ns)
|
||||
ms, err := loadComposeModuleFull(ctx, s, req, ns.ID, "mod1")
|
||||
req.NotNil(ms)
|
||||
|
||||
rr, _, err := store.SearchComposeRecords(ctx, s, ms, types.RecordFilter{})
|
||||
req.NoError(err)
|
||||
req.Len(rr, 2000)
|
||||
|
||||
for i, r := range rr {
|
||||
req.Equal(fmt.Sprintf("r%df1", i+1), r.Values.Get("f1", 0).Value)
|
||||
req.Equal(fmt.Sprintf("r%df2", i+1), r.Values.Get("f2", 0).Value)
|
||||
req.Equal(fmt.Sprintf("r%df3", i+1), r.Values.Get("f3", 0).Value)
|
||||
req.Equal(fmt.Sprintf("r%df4", i+1), r.Values.Get("f4", 0).Value)
|
||||
req.Equal(fmt.Sprintf("r%df5", i+1), r.Values.Get("f5", 0).Value)
|
||||
req.Equal(fmt.Sprintf("r%df6", i+1), r.Values.Get("f6", 0).Value)
|
||||
}
|
||||
|
||||
s.TruncateComposeRecords(ctx, ms)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestDataShaping_fieldTypes(t *testing.T) {
|
||||
var (
|
||||
ctx = auth.SetSuperUserContext(context.Background())
|
||||
|
||||
@@ -0,0 +1,43 @@
|
||||
namespaces:
|
||||
ns1:
|
||||
name: ns1 name
|
||||
|
||||
modules:
|
||||
mod1:
|
||||
fields:
|
||||
f1:
|
||||
label: f1 label
|
||||
kind: String
|
||||
f2:
|
||||
label: f2 label
|
||||
kind: String
|
||||
f3:
|
||||
label: f3 label
|
||||
kind: String
|
||||
f4:
|
||||
label: f4 label
|
||||
kind: String
|
||||
f5:
|
||||
label: f5 label
|
||||
kind: String
|
||||
f6:
|
||||
label: f6 label
|
||||
kind: String
|
||||
|
||||
records:
|
||||
source: mod1.csv
|
||||
key: id
|
||||
mapping:
|
||||
id: /
|
||||
f1:
|
||||
field: f1
|
||||
f2:
|
||||
field: f2
|
||||
f3:
|
||||
field: f3
|
||||
f4:
|
||||
field: f4
|
||||
f5:
|
||||
field: f5
|
||||
f6:
|
||||
field: f6
|
||||
+2001
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,43 @@
|
||||
namespaces:
|
||||
ns1:
|
||||
name: ns1 name
|
||||
|
||||
modules:
|
||||
mod1:
|
||||
fields:
|
||||
f1:
|
||||
label: f1 label
|
||||
kind: String
|
||||
f2:
|
||||
label: f2 label
|
||||
kind: String
|
||||
f3:
|
||||
label: f3 label
|
||||
kind: String
|
||||
f4:
|
||||
label: f4 label
|
||||
kind: String
|
||||
f5:
|
||||
label: f5 label
|
||||
kind: String
|
||||
f6:
|
||||
label: f6 label
|
||||
kind: String
|
||||
|
||||
records:
|
||||
source: mod1.jsonl
|
||||
key: id
|
||||
mapping:
|
||||
id: /
|
||||
f1:
|
||||
field: f1
|
||||
f2:
|
||||
field: f2
|
||||
f3:
|
||||
field: f3
|
||||
f4:
|
||||
field: f4
|
||||
f5:
|
||||
field: f5
|
||||
f6:
|
||||
field: f6
|
||||
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user