Add SQL3 Tests added bulk insert tests (#2324)

* added bulk insert tests

* test idset,stringset,bool in parquet

* bulk insert time coverage
This commit is contained in:
tgruben 2023-03-16 13:02:39 -05:00 committed by GitHub
parent ada48be181
commit c9c63b22b4
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
2 changed files with 135 additions and 21 deletions

View file

@ -266,13 +266,19 @@ func (i *bulkInsertSourceCSVRowIter) Next(ctx context.Context) (types.Row, error
rec, err := i.csvReader.Read()
if err == io.EOF {
return nil, types.ErrNoMoreRows
} else if err != nil {
pe, ok := err.(*csv.ParseError)
if ok {
return nil, sql3.NewErrReadingDatasource(0, 0, i.options.sourceData, fmt.Sprintf("csv parse error on line %d: %s", pe.Line, pe.Error()))
}
return nil, err
}
// err == csv.ParseError is impossible if LazyQuotes is true
// we should uncomment the code below if we ever dispable LazyQuotes
// so there is no need check for any other error
/*
else if err != nil {
pe, ok := err.(*csv.ParseError)
if ok {
return nil, sql3.NewErrReadingDatasource(0, 0, i.options.sourceData, fmt.Sprintf("csv parse error on line %d: %s", pe.Line, pe.Error()))
}
return nil, err
}
*/
// now we do the mapping to the output row
result := make([]interface{}, len(i.options.mapExpressions))
@ -313,8 +319,6 @@ func (i *bulkInsertSourceCSVRowIter) Next(ctx context.Context) (types.Row, error
if err != nil {
if tm, err := time.ParseInLocation(time.RFC3339Nano, evalValue, time.UTC); err == nil {
result[idx] = tm
} else if tm, err := time.ParseInLocation(time.RFC3339, evalValue, time.UTC); err == nil {
result[idx] = tm
} else if tm, err := time.ParseInLocation("2006-01-02", evalValue, time.UTC); err == nil {
result[idx] = tm
} else {
@ -659,8 +663,6 @@ func (i *bulkInsertSourceNDJsonRowIter) Next(ctx context.Context) (types.Row, er
case string:
if tm, err := time.ParseInLocation(time.RFC3339Nano, v, time.UTC); err == nil {
result[idx] = tm
} else if tm, err := time.ParseInLocation(time.RFC3339, v, time.UTC); err == nil {
result[idx] = tm
} else if tm, err := time.ParseInLocation("2006-01-02", v, time.UTC); err == nil {
result[idx] = tm
} else {
@ -1153,8 +1155,6 @@ func (i *bulkInsertSourceParquetRowIter) Next(ctx context.Context) (types.Row, e
} else if stringVal, ok := evalValue.(string); ok {
if tm, err := time.ParseInLocation(time.RFC3339Nano, stringVal, time.UTC); err == nil {
result[idx] = tm
} else if tm, err := time.ParseInLocation(time.RFC3339, stringVal, time.UTC); err == nil {
result[idx] = tm
} else if tm, err := time.ParseInLocation("2006-01-02", stringVal, time.UTC); err == nil {
result[idx] = tm
} else {

View file

@ -1718,6 +1718,16 @@ func TestPlanner_BulkInsert(t *testing.T) {
t.Fatalf("unexpected error: %v", err)
}
})
t.Run("BulkNDJSONFileBadBatchSize", func(t *testing.T) {
_, _, _, err = sql_test.MustQueryRows(t, c.GetNode(0).Server, `bulk insert into j (_id, a, b) map ('id' id, 'a' int, 'b' int) from '/foo/bar' WITH FORMAT 'NDJSON' INPUT 'FILE' BATCHSIZE 0;`)
if err == nil || !strings.Contains(err.Error(), `invalid batch size '0'`) {
t.Fatalf("unexpected error: %v", err)
}
_, _, _, err = sql_test.MustQueryRows(t, c.GetNode(0).Server, `bulk insert into j (_id, a, b) map ('id' id, 'a' int, 'b' int) from '/foo/bar' WITH FORMAT 'NDJSON' INPUT 'FILE' BATCHSIZE 'foo';`)
if err == nil || !strings.Contains(err.Error(), `integer literal expected`) {
t.Fatalf("unexpected error: %v", err)
}
})
t.Run("BulkBadRowsLimit", func(t *testing.T) {
_, _, _, err = sql_test.MustQueryRows(t, c.GetNode(0).Server, `bulk insert into j (_id, a, b) map (0 id, 1 int, 2 int) from '/foo/bar' WITH FORMAT 'CSV' INPUT 'FILE' ROWSLIMIT 'foo';`)
@ -1869,7 +1879,9 @@ func TestPlanner_BulkInsert(t *testing.T) {
x'{ "_id": 1, "id1": 10, "i1": 11, "ids1": [ 3, 4, 5 ], "ss1": [ "foo", "bar" ], "ts1": "2012-11-01T22:08:41+00:00", "s1": "frobny", "b1": true, "d1": 11.34 }
{ "_id": 2, "id1": 10, "i1": 11, "ids1": [ 3, 4, 5 ], "ss1": [ "foo", "bar" ], "ts1": "2012-11-01T22:08:41+00:00", "s1": "frobny", "b1": 0, "d1": 11.34 }
{ "_id": 3, "id1": 10, "i1": 11, "ids1": [ 3, 4, 5 ], "ss1": [ "foo", "bar" ], "ts1": "2012-11-01T22:08:41+00:00", "s1": "frobny", "b1": 1, "d1": 11.34 }
{ "_id": 4, "id1": 10, "i1": 11, "ids1": 9, "ss1": "baz", "ts1": "2012-11-01T22:08:41+00:00", "s1": "frobny", "b1": 1, "d1": 11.34 }'
{ "_id": 4, "id1": 10, "i1": 11, "ids1": 9, "ss1": "baz", "ts1": "2012-11-01T22:08:41+00:00", "s1": "frobny", "b1": 1, "d1": 11.34 }
{ "_id": 5, "id1": 10, "i1": 11, "ids1": 9, "ss1": "baz", "ts1": "2012-11-01", "s1": 42, "b1": 1, "d1": 11.34 }
'
with
format 'NDJSON'
input 'STREAM';`)
@ -2840,6 +2852,11 @@ func simpleParquetMaker(t *testing.T, f *os.File, numRows int64, input []tb) {
sbuild.AppendValues(input[i].Value.([]string), nil)
newChunk := sbuild.NewArray()
chunks = append(chunks, newChunk)
case arrow.FixedWidthTypes.Boolean:
bbuild := array.NewBooleanBuilder(mem)
bbuild.AppendValues(input[i].Value.([]bool), nil)
newChunk := bbuild.NewArray()
chunks = append(chunks, newChunk)
}
}
// make schema
@ -2862,26 +2879,34 @@ func TestPlanner_BulkInsertParquet(t *testing.T) {
t.Run("BulkFromLocalFile", func(t *testing.T) {
// check that can pull parquet file from local file
_, _, _, err := sql_test.MustQueryRows(t, c.GetNode(0).Server, `create table j1 (_id ID, a INT, b DECIMAL(2), c STRING);`)
_, _, _, err := sql_test.MustQueryRows(t, c.GetNode(0).Server, `create table j1 (_id ID, a INT, b DECIMAL(2), c STRING, d STRINGSET, e IDSET, f BOOL, t TIMESTAMP);`)
assert.NoError(t, err)
tmpfile, err := os.CreateTemp("", "BulkParquetFile.parquet")
assert.NoError(t, err)
defer os.Remove(tmpfile.Name())
// create a parquet file with all the example data
simpleParquetMaker(t, tmpfile, 2, []tb{
{Name: "id", Type: arrow.PrimitiveTypes.Int64, Value: []int64{1, 2}},
{Name: "int64V", Type: arrow.PrimitiveTypes.Int64, Value: []int64{42, 7}},
{Name: "float64V", Type: arrow.PrimitiveTypes.Float64, Value: []float64{3.14159, 1.61803}},
{Name: "stringV", Type: arrow.BinaryTypes.String, Value: []string{"pi", "goldenratio"}},
{Name: "id", Type: arrow.PrimitiveTypes.Int64, Value: []int64{1, 2, 3}},
{Name: "int64V", Type: arrow.PrimitiveTypes.Int64, Value: []int64{42, 7, 6}},
{Name: "float64V", Type: arrow.PrimitiveTypes.Float64, Value: []float64{3.14159, 1.61803, 1.41426}},
{Name: "stringV", Type: arrow.BinaryTypes.String, Value: []string{"pi", "goldenratio", "sqr2"}},
{Name: "stringsetV", Type: arrow.BinaryTypes.String, Value: []string{"a1", "a2", "a3"}},
{Name: "idsetV", Type: arrow.PrimitiveTypes.Int64, Value: []int64{10, 20, 30}},
{Name: "boolV", Type: arrow.FixedWidthTypes.Boolean, Value: []bool{true, false, true}},
{Name: "timestampstring", Type: arrow.BinaryTypes.String, Value: []string{"2022-01-28T12:14:04Z", "1970-01-28", "1988-05-30T12:02:00.567999999Z"}},
})
_, _, _, err = sql_test.MustQueryRows(t, c.GetNode(0).Server, fmt.Sprintf(`bulk insert
into j1 (_id,a,b,c )
into j1 (_id,a,b,c,d,e,f,t )
map(
'id' id,
'int64V' INT,
'float64V' DECIMAL(2),
'stringV' STRING)
'stringV' STRING,
'stringsetV' STRINGSET,
'idsetV' IDSET,
'boolV' BOOL,
'timestampstring' TIMESTAMP)
from
'%s'
WITH FORMAT 'PARQUET'
@ -2892,6 +2917,7 @@ func TestPlanner_BulkInsertParquet(t *testing.T) {
if diff := cmp.Diff([][]interface{}{
{int64(1), int64(42), "pi"},
{int64(2), int64(7), "goldenratio"},
{int64(3), int64(6), "sqr2"},
}, results); diff != "" {
t.Fatal(diff)
}
@ -2903,7 +2929,7 @@ func TestPlanner_BulkInsertParquet(t *testing.T) {
}
})
t.Run("BulkFromUrl", func(t *testing.T) {
t.Run("BulkParquetFromUrl", func(t *testing.T) {
// check that can pull parquet file from URL and load
_, _, _, err := sql_test.MustQueryRows(t, c.GetNode(0).Server, `create table j2 (_id ID, a INT, b STRING);`)
@ -2979,6 +3005,94 @@ func TestPlanner_BulkInsertParquet(t *testing.T) {
row := results[0]
assert.Equal(t, row[1], row[2])
})
t.Run("BulkCSVFromUrl", func(t *testing.T) {
// check that can pull parquet file from URL and load
_, _, _, err := sql_test.MustQueryRows(t, c.GetNode(0).Server, `create table j3 (_id ID, a INT, b STRING);`)
assert.NoError(t, err)
tmpfile, err := os.CreateTemp("", "BulkCSVFile.csv")
assert.NoError(t, err)
defer os.Remove(tmpfile.Name())
// create a csv file with all the example data bigger than batch size
var expect [][]interface{}
for i := int64(1); i < 100; i++ {
s := fmt.Sprintf("pi%d", i)
tmpfile.Write([]byte(fmt.Sprintf("%d, %d, \"%s\"\n", i, i+42, s)))
expect = append(expect, []interface{}{i, i + 42, s})
}
mux := http.NewServeMux()
ts := httptest.NewServer(mux)
defer ts.Close()
mux.HandleFunc("/static", func(w http.ResponseWriter, r *http.Request) {
payload, _ := os.ReadFile(tmpfile.Name())
w.Header().Set("Content-Type", "application/octet-stream")
w.WriteHeader(200)
w.Write(payload)
})
_, _, _, err = sql_test.MustQueryRows(t, c.GetNode(0).Server, fmt.Sprintf(`bulk insert
into j3 (_id,a,b )
map(
0 id,
1 INT,
2 STRING)
from
'%s'
WITH FORMAT 'CSV'
BATCHSIZE 60
INPUT 'URL';`, ts.URL+"/static"))
assert.NoError(t, err)
results, _, _, err := sql_test.MustQueryRows(t, c.GetNode(0).Server, `select _id, a,b from j3`)
assert.NoError(t, err)
if diff := cmp.Diff(expect, results); diff != "" {
t.Fatal(diff)
}
})
t.Run("BulkNDJSONFromUrl", func(t *testing.T) {
// check that can pull parquet file from URL and load
_, _, _, err := sql_test.MustQueryRows(t, c.GetNode(0).Server, `create table j4 (_id ID, a INT, b STRING, c TIMESTAMP);`)
assert.NoError(t, err)
tmpfile, err := os.CreateTemp("", "BulkJSONFile.json")
assert.NoError(t, err)
defer os.Remove(tmpfile.Name())
// create a csv file with all the example data
tmpfile.Write([]byte(`{ "id":1, "intv":42, "stringv":"pi" , "ts": 1678985262}`))
mux := http.NewServeMux()
ts := httptest.NewServer(mux)
defer ts.Close()
mux.HandleFunc("/static", func(w http.ResponseWriter, r *http.Request) {
payload, _ := os.ReadFile(tmpfile.Name())
w.Header().Set("Content-Type", "application/octet-stream")
w.WriteHeader(200)
w.Write(payload)
})
_, _, _, err = sql_test.MustQueryRows(t, c.GetNode(0).Server, fmt.Sprintf(`bulk insert
into j4 (_id,a,b,c )
map(
'$.id' id,
'$.intv' INT,
'$.stringv' STRING,
'$.ts' TIMESTAMP)
from
'%s'
WITH FORMAT 'NDJSON'
INPUT 'URL';`, ts.URL+"/static"))
assert.NoError(t, err)
results, _, _, err := sql_test.MustQueryRows(t, c.GetNode(0).Server, `select _id, a,b from j4`)
assert.NoError(t, err)
if diff := cmp.Diff([][]interface{}{
{int64(1), int64(42), "pi"},
}, results); diff != "" {
t.Fatal(diff)
}
})
}
func TestPlanner_BulkInsert_FP1916(t *testing.T) {