// Licensed to the LF AI & Data foundation under one // or more contributor license agreements. See the NOTICE file // distributed with this work for additional information // regarding copyright ownership. The ASF licenses this file // to you under the Apache License, Version 2.0 (the // "License"); you may not use this file except in compliance // with the License. You may obtain a copy of the License at // // http://www.apache.org/licenses/LICENSE-2.0 // // Unless required by applicable law or agreed to in writing, software // distributed under the License is distributed on an "AS IS" BASIS, // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. // See the License for the specific language governing permissions and // limitations under the License. package csv import ( "context" "encoding/csv" "fmt" "io" "math/rand" "os" "strings" "testing" "github.com/stretchr/testify/mock" "github.com/stretchr/testify/suite" "github.com/milvus-io/milvus-proto/go-api/v3/commonpb" "github.com/milvus-io/milvus-proto/go-api/v3/schemapb" "github.com/milvus-io/milvus/internal/mocks" "github.com/milvus-io/milvus/internal/storage" importcommon "github.com/milvus-io/milvus/internal/util/importutilv2/common" "github.com/milvus-io/milvus/internal/util/testutil" "github.com/milvus-io/milvus/pkg/v3/common" "github.com/milvus-io/milvus/pkg/v3/objectstorage" "github.com/milvus-io/milvus/pkg/v3/util/merr" "github.com/milvus-io/milvus/pkg/v3/util/paramtable" ) func init() { paramtable.Init() } type ReaderSuite struct { suite.Suite numRows int pkDataType schemapb.DataType vecDataType schemapb.DataType } func (suite *ReaderSuite) SetupTest() { suite.numRows = 100 suite.pkDataType = schemapb.DataType_Int64 suite.vecDataType = schemapb.DataType_FloatVector } func (suite *ReaderSuite) run(dataType schemapb.DataType, elemType schemapb.DataType, nullable bool, nullPercent int) { schema := &schemapb.CollectionSchema{ Fields: []*schemapb.FieldSchema{ { FieldID: 100, Name: "pk", IsPrimaryKey: true, DataType: suite.pkDataType, TypeParams: []*commonpb.KeyValuePair{ { Key: common.MaxLengthKey, Value: "128", }, }, }, { FieldID: 101, Name: "vec", DataType: suite.vecDataType, TypeParams: []*commonpb.KeyValuePair{ { Key: common.DimKey, Value: "8", }, }, Nullable: nullable, }, { FieldID: 102, Name: dataType.String(), DataType: dataType, ElementType: elemType, TypeParams: []*commonpb.KeyValuePair{ { Key: common.MaxLengthKey, Value: "128", }, { Key: common.MaxCapacityKey, Value: "256", }, }, Nullable: nullable, }, }, } if dataType == schemapb.DataType_VarChar { // Add a BM25 function if data type is VarChar schema.Fields = append(schema.Fields, &schemapb.FieldSchema{ FieldID: 103, Name: "sparse", DataType: schemapb.DataType_SparseFloatVector, IsFunctionOutput: true, }) schema.Functions = append(schema.Functions, &schemapb.FunctionSchema{ Id: 1000, Name: "bm25", Type: schemapb.FunctionType_BM25, InputFieldIds: []int64{102}, InputFieldNames: []string{dataType.String()}, OutputFieldIds: []int64{103}, OutputFieldNames: []string{"sparse"}, }) } // config // csv separator sep := ',' // csv writer write null value as empty string nullkey := "" // generate csv data insertData, err := testutil.CreateInsertData(schema, suite.numRows, nullPercent) suite.NoError(err) csvData, err := testutil.CreateInsertDataForCSV(schema, insertData, nullkey) suite.NoError(err) // write to csv file filePath := fmt.Sprintf("/tmp/test_%d_reader.csv", rand.Int()) defer os.Remove(filePath) wf, err := os.OpenFile(filePath, os.O_RDWR|os.O_CREATE, 0o666) suite.NoError(err) writer := csv.NewWriter(wf) writer.Comma = sep err = writer.WriteAll(csvData) suite.NoError(err) // read from csv file ctx := context.Background() f := storage.NewChunkManagerFactory("local", objectstorage.RootPath(suite.T().TempDir()+"/")) cm, err := f.NewPersistentStorageChunkManager(ctx) suite.NoError(err) // check reader separate fields by '\t' wrongSep := '\t' _, err = NewReader(ctx, cm, schema, filePath, 64*1024*1024, wrongSep, nullkey) suite.Error(err) suite.Contains(err.Error(), "value of field is missed:") // check data reader, err := NewReader(ctx, cm, schema, filePath, 64*1024*1024, sep, nullkey) suite.NoError(err) checkFn := func(actualInsertData *storage.InsertData, offsetBegin, expectRows int) { expectInsertData := insertData for fieldID, data := range actualInsertData.Data { suite.Equal(expectRows, data.RowNum()) for i := 0; i < expectRows; i++ { expect := expectInsertData.Data[fieldID].GetRow(i + offsetBegin) actual := data.GetRow(i) suite.Equal(expect, actual) } } } res, err := reader.Read() suite.NoError(err) checkFn(res, 0, suite.numRows) } func (suite *ReaderSuite) TestReadScalarFields() { elementTypes := []schemapb.DataType{ schemapb.DataType_Bool, schemapb.DataType_Int8, schemapb.DataType_Int16, schemapb.DataType_Int32, schemapb.DataType_Int64, schemapb.DataType_Float, schemapb.DataType_Double, schemapb.DataType_String, } scalarTypes := append(elementTypes, []schemapb.DataType{schemapb.DataType_VarChar, schemapb.DataType_JSON, schemapb.DataType_Array}...) for _, dataType := range scalarTypes { if dataType == schemapb.DataType_Array { for _, elementType := range elementTypes { suite.run(dataType, elementType, false, 0) for _, nullPercent := range []int{0, 50, 100} { suite.run(dataType, elementType, true, nullPercent) } } } else { suite.run(dataType, schemapb.DataType_None, false, 0) for _, nullPercent := range []int{0, 50, 100} { suite.run(dataType, schemapb.DataType_None, true, nullPercent) } } } } func (suite *ReaderSuite) TestStringPK() { suite.pkDataType = schemapb.DataType_VarChar suite.run(schemapb.DataType_Int32, schemapb.DataType_None, false, 0) for _, nullPercent := range []int{0, 50, 100} { suite.run(schemapb.DataType_Int32, schemapb.DataType_None, true, nullPercent) } } func (suite *ReaderSuite) TestVector() { dataTypes := []schemapb.DataType{ schemapb.DataType_BinaryVector, schemapb.DataType_FloatVector, schemapb.DataType_Float16Vector, schemapb.DataType_BFloat16Vector, schemapb.DataType_SparseFloatVector, schemapb.DataType_Int8Vector, } for _, dataType := range dataTypes { suite.vecDataType = dataType suite.run(schemapb.DataType_Int32, schemapb.DataType_None, false, 0) for _, nullPercent := range []int{0, 50, 100} { suite.run(schemapb.DataType_Int32, schemapb.DataType_None, true, nullPercent) } } } func (suite *ReaderSuite) TestError() { testNewReaderErr := func(ioErr error, schema *schemapb.CollectionSchema, content string, bufferSize int) { cm := mocks.NewChunkManager(suite.T()) cm.EXPECT().Reader(mock.Anything, mock.Anything).RunAndReturn(func(ctx context.Context, s string) (storage.FileReader, error) { if ioErr != nil { return nil, ioErr } else { r := importcommon.NewMockReader(content) return r, nil } }) cm.EXPECT().Size(mock.Anything, "dummy path").Return(int64(len(content)), nil).Maybe() reader, err := NewReader(context.Background(), cm, schema, "dummy path", bufferSize, ',', "") suite.Error(err) suite.Nil(reader) } testNewReaderErr(merr.WrapErrImportFailed("error"), nil, "", 1) schema := &schemapb.CollectionSchema{ Fields: []*schemapb.FieldSchema{ { FieldID: 100, Name: "pk", IsPrimaryKey: true, DataType: schemapb.DataType_Int64, }, { FieldID: 101, Name: "int32", DataType: schemapb.DataType_Int32, }, }, } testNewReaderErr(nil, schema, "", -1) testNewReaderErr(nil, schema, "", 1) testReadErr := func(schema *schemapb.CollectionSchema, content string) { cm := mocks.NewChunkManager(suite.T()) cm.EXPECT().Reader(mock.Anything, mock.Anything).RunAndReturn(func(ctx context.Context, s string) (storage.FileReader, error) { r := importcommon.NewMockReader(content) return r, nil }) cm.EXPECT().Size(mock.Anything, mock.Anything).Return(128, nil) reader, err := NewReader(context.Background(), cm, schema, "dummy path", 1024, ',', "") suite.NoError(err) suite.NotNil(reader) size, err := reader.Size() suite.NoError(err) suite.Equal(int64(128), size) _, err = reader.Read() suite.Error(err) } content := "pk,int32\n1,A" testReadErr(schema, content) } func (suite *ReaderSuite) TestReadLoop() { schema := &schemapb.CollectionSchema{ Fields: []*schemapb.FieldSchema{ { FieldID: 100, Name: "pk", IsPrimaryKey: true, DataType: schemapb.DataType_Int64, }, { FieldID: 101, Name: "float", DataType: schemapb.DataType_Float, }, }, } content := "pk,float\n1,0.1\n2,0.2" cm := mocks.NewChunkManager(suite.T()) cm.EXPECT().Reader(mock.Anything, mock.Anything).RunAndReturn(func(ctx context.Context, s string) (storage.FileReader, error) { r := importcommon.NewMockReader(content) return r, nil }) cm.EXPECT().Size(mock.Anything, mock.Anything).Return(128, nil) reader, err := NewReader(context.Background(), cm, schema, "dummy path", 1, ',', "") suite.NoError(err) suite.NotNil(reader) defer reader.Close() size, err := reader.Size() suite.NoError(err) suite.Equal(int64(128), size) size2, err := reader.Size() // size is cached suite.NoError(err) suite.Equal(size, size2) reader.count = 1 data, err := reader.Read() suite.NoError(err) suite.Equal(1, data.GetRowNum()) data, err = reader.Read() suite.NoError(err) suite.Equal(1, data.GetRowNum()) data, err = reader.Read() suite.EqualError(io.EOF, err.Error()) suite.Nil(data) } func (suite *ReaderSuite) TestReadRecoversFromPrematureEOF() { schema := &schemapb.CollectionSchema{ Fields: []*schemapb.FieldSchema{ { FieldID: 100, Name: "pk", IsPrimaryKey: true, DataType: schemapb.DataType_Int64, }, { FieldID: 101, Name: "int32", DataType: schemapb.DataType_Int32, }, }, } var contentBuilder strings.Builder contentBuilder.WriteString("pk,int32\n") for i := 1; i <= 100; i++ { fmt.Fprintf(&contentBuilder, "%d,%d\n", i, i*10) } content := contentBuilder.String() openCount := 0 cm := mocks.NewChunkManager(suite.T()) cm.EXPECT().Reader(mock.Anything, "dummy path").RunAndReturn(func(ctx context.Context, path string) (storage.FileReader, error) { openCount++ if openCount == 1 { return importcommon.NewPrematureEOFReader(content, 18), nil } return importcommon.NewMockReader(content), nil }) reader, err := NewReader(context.Background(), cm, schema, "dummy path", 1024, ',', "") suite.NoError(err) defer reader.Close() data, err := reader.Read() suite.NoError(err) suite.Equal(100, data.GetRowNum()) suite.Equal(int64(1), data.Data[100].GetRow(0)) suite.Equal(int32(1000), data.Data[101].GetRow(99)) suite.GreaterOrEqual(openCount, 2) } func TestCsvReader(t *testing.T) { suite.Run(t, new(ReaderSuite)) } func (suite *ReaderSuite) TestAllowInsertAutoID_KeepUserPK() { schema := &schemapb.CollectionSchema{ Fields: []*schemapb.FieldSchema{ { FieldID: 100, Name: "pk", IsPrimaryKey: true, DataType: schemapb.DataType_Int64, AutoID: true, }, { FieldID: 101, Name: "vec", DataType: schemapb.DataType_FloatVector, TypeParams: []*commonpb.KeyValuePair{{Key: common.DimKey, Value: "8"}}, }, }, } // prepare csv content with header including pk content := "pk,vec\n" // one row is enough; vec must be a quoted JSON array to avoid CSV splitting content += "1,\"[0,1,2,3,4,5,6,7]\"\n" // allow_insert_autoid=false, providing PK in header should error { cm := mocks.NewChunkManager(suite.T()) cm.EXPECT().Reader(mock.Anything, mock.Anything).RunAndReturn(func(ctx context.Context, s string) (storage.FileReader, error) { r := importcommon.NewMockReader(content) return r, nil }) _, err := NewReader(context.Background(), cm, schema, "dummy path", 1024, ',', "") suite.Error(err) suite.Contains(err.Error(), "is auto-generated, no need to provide") } // allow_insert_autoid=true, providing PK in header should be allowed { schema.Properties = []*commonpb.KeyValuePair{{Key: common.AllowInsertAutoIDKey, Value: "true"}} cm := mocks.NewChunkManager(suite.T()) cm.EXPECT().Reader(mock.Anything, mock.Anything).RunAndReturn(func(ctx context.Context, s string) (storage.FileReader, error) { r := importcommon.NewMockReader(content) return r, nil }) reader, err := NewReader(context.Background(), cm, schema, "dummy path", 1024, ',', "") suite.NoError(err) // call Read once to ensure parsing proceeds _, err = reader.Read() suite.NoError(err) } } func (suite *ReaderSuite) TestFunctionOutputField() { // Test 1: BM25 function output NOT in CSV should be excluded from result, // but nullable non-function-output field not in CSV should be retained (filled with nil). { schema := &schemapb.CollectionSchema{ Fields: []*schemapb.FieldSchema{ {FieldID: 100, Name: "pk", IsPrimaryKey: true, DataType: schemapb.DataType_Int64}, { FieldID: 101, Name: "text", DataType: schemapb.DataType_VarChar, TypeParams: []*commonpb.KeyValuePair{{Key: common.MaxLengthKey, Value: "128"}}, }, {FieldID: 102, Name: "sparse", DataType: schemapb.DataType_SparseFloatVector, IsFunctionOutput: true}, {FieldID: 103, Name: "score", DataType: schemapb.DataType_Float, Nullable: true}, }, Functions: []*schemapb.FunctionSchema{ { Id: 1000, Name: "bm25", Type: schemapb.FunctionType_BM25, InputFieldIds: []int64{101}, InputFieldNames: []string{"text"}, OutputFieldIds: []int64{102}, OutputFieldNames: []string{"sparse"}, }, }, } // CSV does NOT provide "sparse" (function output) or "score" (nullable) content := "pk,text\n1,hello\n2,world" cm := mocks.NewChunkManager(suite.T()) cm.EXPECT().Reader(mock.Anything, mock.Anything).RunAndReturn(func(ctx context.Context, s string) (storage.FileReader, error) { return importcommon.NewMockReader(content), nil }) reader, err := NewReader(context.Background(), cm, schema, "dummy path", 1024, ',', "") suite.NoError(err) res, err := reader.Read() suite.NoError(err) suite.Equal(2, res.GetRowNum()) // BM25 function output field should be removed (not populated from CSV) _, hasSparse := res.Data[102] suite.False(hasSparse) // Nullable non-function-output field should be retained (filled with nil, not removed) scoreData, hasScore := res.Data[103] suite.True(hasScore) suite.Equal(2, scoreData.RowNum()) } // Test 2: Non-BM25 function output in CSV with property enabled should be included { schema := &schemapb.CollectionSchema{ Fields: []*schemapb.FieldSchema{ {FieldID: 100, Name: "pk", IsPrimaryKey: true, DataType: schemapb.DataType_Int64}, { FieldID: 101, Name: "text", DataType: schemapb.DataType_VarChar, TypeParams: []*commonpb.KeyValuePair{{Key: common.MaxLengthKey, Value: "128"}}, }, { FieldID: 102, Name: "embedding", DataType: schemapb.DataType_FloatVector, IsFunctionOutput: true, TypeParams: []*commonpb.KeyValuePair{{Key: common.DimKey, Value: "2"}}, }, }, Functions: []*schemapb.FunctionSchema{ { Id: 1000, Name: "text_embedding", Type: schemapb.FunctionType_TextEmbedding, InputFieldIds: []int64{101}, InputFieldNames: []string{"text"}, OutputFieldIds: []int64{102}, OutputFieldNames: []string{"embedding"}, }, }, Properties: []*commonpb.KeyValuePair{ {Key: common.CollectionAllowInsertNonBM25FunctionOutputs, Value: "true"}, }, } content := "pk,text,embedding\n1,hello,\"[0.1, 0.2]\"\n2,world,\"[0.3, 0.4]\"" cm := mocks.NewChunkManager(suite.T()) cm.EXPECT().Reader(mock.Anything, mock.Anything).RunAndReturn(func(ctx context.Context, s string) (storage.FileReader, error) { return importcommon.NewMockReader(content), nil }) reader, err := NewReader(context.Background(), cm, schema, "dummy path", 1024, ',', "") suite.NoError(err) res, err := reader.Read() suite.NoError(err) suite.Equal(2, res.GetRowNum()) // embedding field should be present with correct data suite.NotNil(res.Data[102]) suite.Equal(2, res.Data[102].RowNum()) } }