enhance: [3.0] pass external segment row target from DataCoord (#51364)

issue: #45881
pr: #51352

## What this PR does

Backports the external refresh segment row-target propagation to 3.0.

- DataCoord snapshots `dataNode.externalCollection.targetRowsPerSegment`
  into each refresh task request.
- DataNode uses the transmitted value for fragment splitting and segment
  balancing.
- Requests from older DataCoord instances fall back to the DataNode
  configuration before validation.
- The existing configuration remains the source of the value.

## Branch notes

This is a clean cherry-pick of the master change. No 3.0-specific code
changes or conflict resolution were required.

## Validation

- `make generated-proto-without-cpp`
- `make build-cpp`
- `make build-go`
- `make lint-fix`
- Targeted race-enabled tests for `internal/datacoord` and
  `internal/datanode/external` with the required dynamic/test tags,
  gcflags, and ldflags
- `git diff --check`

Signed-off-by: Wei Liu <wei.liu@zilliz.com>
This commit is contained in:
wei liu
2026-07-17 10:52:40 +08:00
committed by GitHub
parent 3f52db0cf3
commit c164ffd8ae
6 changed files with 1428 additions and 1341 deletions
@@ -688,6 +688,7 @@ func (t *refreshExternalCollectionTask) CreateTaskOnWorker(nodeID int64, cluster
ExploreManifestPath: t.GetExploreManifestPath(),
FileIndexBegin: t.GetFileIndexBegin(),
FileIndexEnd: t.GetFileIndexEnd(),
TargetRowsPerSegment: paramtable.Get().DataNodeCfg.ExternalCollectionTargetRowsPerSegment.GetAsInt64(),
}
// Submit task to worker via unified task system
@@ -865,6 +865,10 @@ func TestRefreshExternalCollectionTask_CreateTaskOnWorker(t *testing.T) {
})
t.Run("success", func(t *testing.T) {
const targetRowsPerSegmentKey = "dataNode.externalCollection.targetRowsPerSegment"
paramtable.Get().Save(targetRowsPerSegmentKey, "12345")
defer paramtable.Get().Reset(targetRowsPerSegmentKey)
catalog := &stubCatalog{}
refreshMeta, err := newExternalCollectionRefreshMeta(context.Background(), catalog)
assert.NoError(t, err)
@@ -912,6 +916,7 @@ func TestRefreshExternalCollectionTask_CreateTaskOnWorker(t *testing.T) {
assert.Equal(t, indexpb.JobState_JobStateInProgress, metaTask.GetState())
assert.NotNil(t, cluster.refreshReq)
assert.Equal(t, int64(10), cluster.refreshReq.GetPartitionID())
assert.Equal(t, int64(12345), cluster.refreshReq.GetTargetRowsPerSegment())
})
}
+8 -5
View File
@@ -187,6 +187,12 @@ func (t *RefreshExternalCollectionTask) PreExecute(ctx context.Context) error {
if err != nil {
return merr.Wrap(err, "failed to parse external spec")
}
if t.req.GetTargetRowsPerSegment() == 0 {
t.req.TargetRowsPerSegment = paramtable.Get().DataNodeCfg.ExternalCollectionTargetRowsPerSegment.GetAsInt64()
}
if t.req.GetTargetRowsPerSegment() <= 0 {
return merr.WrapErrParameterInvalidMsg("target rows per segment must be positive")
}
t.parsedSpec = spec
schema := proto.Clone(t.req.GetSchema()).(*schemapb.CollectionSchema)
if schema.GetExternalSource() == "" {
@@ -256,8 +262,6 @@ func (t *RefreshExternalCollectionTask) fetchFragmentsFromExternalSource(ctx con
mlog.Int64("fileIndexBegin", t.req.GetFileIndexBegin()),
mlog.Int64("fileIndexEnd", t.req.GetFileIndexEnd()))
targetRowsPerSegment := paramtable.Get().DataNodeCfg.ExternalCollectionTargetRowsPerSegment.GetAsInt64()
return packed.FetchFragmentsFromExternalSourceWithRange(
ctx,
t.parsedSpec.Format,
@@ -270,7 +274,7 @@ func (t *RefreshExternalCollectionTask) fetchFragmentsFromExternalSource(ctx con
packed.ExternalFetchOptions{
CollectionID: t.req.GetCollectionID(),
ExternalSpec: t.req.GetExternalSpec(),
RowLimit: targetRowsPerSegment,
RowLimit: t.req.GetTargetRowsPerSegment(),
},
)
}
@@ -720,8 +724,7 @@ func (t *RefreshExternalCollectionTask) balanceFragmentsToSegments(ctx context.C
}
}
} else {
// Get target rows per segment from configuration
targetRowsPerSegment := paramtable.Get().DataNodeCfg.ExternalCollectionTargetRowsPerSegment.GetAsInt64()
targetRowsPerSegment := t.req.GetTargetRowsPerSegment()
if totalRows < targetRowsPerSegment {
targetRowsPerSegment = totalRows
}
+146 -80
View File
@@ -51,6 +51,16 @@ type RefreshExternalCollectionTaskSuite struct {
taskID int64
}
func (s *RefreshExternalCollectionTaskSuite) newTask(
ctx context.Context,
req *datapb.RefreshExternalCollectionTaskRequest,
) *RefreshExternalCollectionTask {
if req != nil && req.GetTargetRowsPerSegment() == 0 {
req.TargetRowsPerSegment = 1_000_000
}
return NewRefreshExternalCollectionTask(ctx, req)
}
type fakeRecordReader struct {
records []storage.Record
errs []error
@@ -90,7 +100,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestNewRefreshExternalCollectionTas
ExternalSpec: "test_spec",
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
s.NotNil(task)
s.Equal(s.collectionID, task.req.GetCollectionID())
@@ -125,7 +135,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestTaskLifecycle() {
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
// Test OnEnqueue
err := task.OnEnqueue(ctx)
@@ -166,7 +176,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestPreExecuteClonesSchemaBeforeFil
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
err := task.PreExecute(ctx)
s.NoError(err)
@@ -178,6 +188,48 @@ func (s *RefreshExternalCollectionTaskSuite) TestPreExecuteClonesSchemaBeforeFil
s.Equal([]string{"id"}, task.columns)
}
func (s *RefreshExternalCollectionTaskSuite) TestPreExecuteFallsBackToConfigWhenTargetRowsPerSegmentMissing() {
paramtable.Init()
const targetRowsPerSegmentKey = "dataNode.externalCollection.targetRowsPerSegment"
paramtable.Get().Save(targetRowsPerSegmentKey, "12345")
defer paramtable.Get().Reset(targetRowsPerSegmentKey)
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
req := &datapb.RefreshExternalCollectionTaskRequest{
CollectionID: s.collectionID,
TaskID: s.taskID,
ExternalSource: "test_source",
ExternalSpec: `{"format":"parquet"}`,
Schema: &schemapb.CollectionSchema{},
StorageConfig: &indexpb.StorageConfig{StorageType: "local"},
}
task := NewRefreshExternalCollectionTask(ctx, req)
err := task.PreExecute(ctx)
s.NoError(err)
s.Equal(int64(12345), req.GetTargetRowsPerSegment())
}
func (s *RefreshExternalCollectionTaskSuite) TestPreExecuteRejectsNegativeTargetRowsPerSegment() {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
req := &datapb.RefreshExternalCollectionTaskRequest{
CollectionID: s.collectionID,
TaskID: s.taskID,
ExternalSource: "test_source",
ExternalSpec: `{"format":"parquet"}`,
Schema: &schemapb.CollectionSchema{},
StorageConfig: &indexpb.StorageConfig{StorageType: "local"},
TargetRowsPerSegment: -1,
}
task := NewRefreshExternalCollectionTask(ctx, req)
err := task.PreExecute(ctx)
s.Error(err)
s.Contains(err.Error(), "target rows per segment must be positive")
}
func (s *RefreshExternalCollectionTaskSuite) TestPreExecuteWithNilRequest() {
ctx, cancel := context.WithCancel(context.Background()) //nolint:gosec // cancel is deferred below
defer cancel()
@@ -198,7 +250,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestSetAndGetState() {
TaskID: s.taskID,
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.SetState(indexpb.JobState_JobStateInProgress, "")
s.Equal(indexpb.JobState_JobStateInProgress, task.GetState())
@@ -216,7 +268,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestReset() {
TaskID: s.taskID,
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.Reset()
s.Nil(task.ctx)
@@ -233,7 +285,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Empt
TaskID: s.taskID,
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
result, err := task.balanceFragmentsToSegments(context.Background(), []packed.Fragment{})
s.NoError(err)
s.Nil(result)
@@ -251,7 +303,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Zero
CollectionID: s.collectionID,
TaskID: s.taskID,
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
fragments := []packed.Fragment{
{FragmentID: 1, RowCount: 0, FilePath: "s3://bucket/zero.parquet"},
@@ -285,7 +337,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Sing
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -332,6 +384,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Mult
req := &datapb.RefreshExternalCollectionTaskRequest{
CollectionID: s.collectionID,
TaskID: s.taskID,
TargetRowsPerSegment: 500000,
PreAllocatedSegmentIds: &datapb.IDRange{Begin: 100, End: 200},
StorageConfig: &indexpb.StorageConfig{StorageType: "local"},
ExternalSource: "s3://bucket/data/",
@@ -343,7 +396,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Mult
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -368,6 +421,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Mult
result, err := task.balanceFragmentsToSegments(context.Background(), fragments)
s.NoError(err)
s.Len(result, 4)
// Verify all segments have StorageVersion=V3 and fake binlogs
for i, seg := range result {
@@ -413,7 +467,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestPreExecuteContextCanceled() {
TaskID: s.taskID,
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
cancel()
err := task.PreExecute(ctx)
@@ -428,7 +482,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestExecuteContextCanceled() {
TaskID: s.taskID,
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
cancel()
err := task.Execute(ctx)
@@ -443,7 +497,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegmentsConte
TaskID: s.taskID,
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
cancel()
result, err := task.balanceFragmentsToSegments(ctx, []packed.Fragment{{FragmentID: 1, RowCount: 10}})
@@ -463,7 +517,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestOrganizeSegments_AllFragmentsEx
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
// Simulate current segment fragments mapping (use FilePath as identifier)
currentSegmentFragments := packed.SegmentFragments{
@@ -526,7 +580,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestOrganizeSegments_SameFragmentsM
}, "parquet"),
}},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
task.columns = []string{"id", "vec", "text", "score"}
@@ -625,7 +679,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestOrganizeSegments_PatchedSegment
}, "parquet"),
}},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
task.columns = []string{"id", "text", "score"}
@@ -686,7 +740,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestOrganizeSegments_SameFragmentsA
}, "parquet"),
}},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
current := packed.SegmentFragments{
10: []packed.Fragment{{FilePath: "s3://bucket/data/a.parquet", StartRow: 0, EndRow: 100, RowCount: 100}},
}
@@ -717,7 +771,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestOrganizeSegments_PartialFragmen
PreAllocatedSegmentIds: &datapb.IDRange{Begin: 100, End: 200},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
// S1 has file1, S2 has file2, S3 has file3
currentSegmentFragments := packed.SegmentFragments{
@@ -751,7 +805,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestOrganizeSegments_NewFragmentsUs
{ID: 1, CollectionID: s.collectionID, NumOfRows: 1000},
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
currentSegmentFragments := packed.SegmentFragments{
1: []packed.Fragment{{FragmentID: 101, FilePath: "/data/file1.parquet", StartRow: 0, EndRow: 1000, RowCount: 1000}},
@@ -806,7 +860,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestOrganizeSegments_RewritesSegmen
},
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
currentSegmentFragments := packed.SegmentFragments{
1: []packed.Fragment{{FragmentID: 101, FilePath: "/data/file1.parquet", StartRow: 0, EndRow: 1000, RowCount: 1000}},
@@ -869,7 +923,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestOrganizeSegments_KeepsSegmentWi
},
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
currentSegmentFragments := packed.SegmentFragments{
1: []packed.Fragment{{FragmentID: 101, FilePath: "/data/file1.parquet", StartRow: 0, EndRow: 1000, RowCount: 1000}},
@@ -922,7 +976,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestOrganizeSegments_FunctionOutput
},
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
currentSegmentFragments := packed.SegmentFragments{
1: []packed.Fragment{{FragmentID: 101, FilePath: "/data/file1.parquet", StartRow: 0, EndRow: 1000, RowCount: 1000}},
@@ -1125,7 +1179,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBuildFunctionExecutionSchemaBra
}
func (s *RefreshExternalCollectionTaskSuite) TestSegmentHasFunctionOutputColumnsBranches() {
task := NewRefreshExternalCollectionTask(context.Background(), &datapb.RefreshExternalCollectionTaskRequest{})
task := s.newTask(context.Background(), &datapb.RefreshExternalCollectionTaskRequest{})
hasColumns, err := task.segmentHasFunctionOutputColumns(&datapb.SegmentInfo{}, nil)
s.NoError(err)
@@ -1165,7 +1219,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestOrganizeSegmentsFunctionOutputC
},
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
result, err := task.organizeSegments(ctx, nil, nil)
s.Error(err)
@@ -1183,7 +1237,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestOrganizeSegmentsBalanceError()
CollectionID: s.collectionID,
TaskID: s.taskID,
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
mockBalance := mockey.Mock(mockey.GetMethod(task, "balanceFragmentsToSegments")).
Return(nil, fmt.Errorf("balance failed")).Build()
defer mockBalance.UnPatch()
@@ -1234,7 +1288,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestOrganizeSegmentsContextCanceled
TaskID: s.taskID,
CurrentSegments: tc.currentSegments,
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
var calls int
mockEnsure := mockey.Mock(ensureContext).To(func(ctx context.Context) error {
@@ -1359,7 +1413,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestCreateManifestForSegment() {
PartitionID: 2000,
StorageConfig: &indexpb.StorageConfig{RootPath: "files", StorageType: "local"},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
task.columns = []string{"col1"}
@@ -1416,7 +1470,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestCreateManifestWithFunctionsUses
},
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
var gotBasePath string
@@ -2167,7 +2221,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestTaskAccessorsAndCloneHelpers()
s.Error(ensureContext(ctx))
taskCtx := context.Background()
task := NewRefreshExternalCollectionTask(taskCtx, &datapb.RefreshExternalCollectionTaskRequest{TaskID: 10})
task := s.newTask(taskCtx, &datapb.RefreshExternalCollectionTaskRequest{TaskID: 10})
task.updatedSegments = []*datapb.SegmentInfo{{ID: 100}}
s.Equal(taskCtx, task.Ctx())
s.Equal(task.updatedSegments, task.GetUpdatedSegments())
@@ -2212,7 +2266,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestExecuteWithMockedSteps() {
FileIndexBegin: 0,
FileIndexEnd: 1,
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
task.columns = []string{"text"}
fragments := []packed.Fragment{{FragmentID: 1, FilePath: "file", StartRow: 0, EndRow: 10, RowCount: 10}}
@@ -2245,10 +2299,10 @@ func (s *RefreshExternalCollectionTaskSuite) TestExecuteErrorPathsWithMockedStep
TaskID: s.taskID,
PreAllocatedSegmentIds: &datapb.IDRange{Begin: 3000, End: 4000},
}
return NewRefreshExternalCollectionTask(ctx, req)
return s.newTask(ctx, req)
}
task := NewRefreshExternalCollectionTask(ctx, &datapb.RefreshExternalCollectionTaskRequest{TaskID: s.taskID})
task := s.newTask(ctx, &datapb.RefreshExternalCollectionTaskRequest{TaskID: s.taskID})
s.Error(task.Execute(ctx))
task = newTask()
@@ -2283,7 +2337,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestExecuteErrorPathsWithMockedStep
func (s *RefreshExternalCollectionTaskSuite) TestPostExecuteAndOrganizeCanceled() {
ctx, cancel := context.WithCancel(context.Background())
cancel()
task := NewRefreshExternalCollectionTask(ctx, &datapb.RefreshExternalCollectionTaskRequest{
task := s.newTask(ctx, &datapb.RefreshExternalCollectionTaskRequest{
TaskID: s.taskID,
CollectionID: s.collectionID,
})
@@ -2304,7 +2358,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBuildCurrentSegmentFragments()
StorageConfig: &indexpb.StorageConfig{RootPath: "files", StorageType: "local"},
CurrentSegments: []*datapb.SegmentInfo{{ID: 1}},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.columns = []string{"text_col"}
expected := packed.SegmentFragments{1: []packed.Fragment{{FragmentID: 1}}}
var gotColumns []string
@@ -2347,7 +2401,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegmentsWithF
},
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -2414,7 +2468,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestPreAllocatedSegmentIDs() {
CurrentSegments: []*datapb.SegmentInfo{},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
// Before Execute(), pre-allocated fields should not be initialized
s.Nil(task.preallocatedIDRange)
@@ -2449,7 +2503,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestPreAllocatedIDAllocation() {
CurrentSegments: []*datapb.SegmentInfo{},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
// Manually initialize (simulating Execute)
task.preallocatedIDRange = idRange
@@ -2478,7 +2532,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestMissingPreAllocatedIDs() {
// PreAllocatedSegmentIds is nil
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
// Execute should fail because pre-allocated IDs are missing
err := task.Execute(ctx)
@@ -2496,7 +2550,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestPreExecute_NilSchema() {
StorageConfig: &indexpb.StorageConfig{StorageType: "local"},
// Schema is nil
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
err := task.PreExecute(ctx)
s.Error(err)
s.Contains(err.Error(), "schema is nil")
@@ -2512,7 +2566,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestPreExecute_NilStorageConfig() {
Schema: &schemapb.CollectionSchema{},
// StorageConfig is nil
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
err := task.PreExecute(ctx)
s.Error(err)
s.Contains(err.Error(), "storage config is nil")
@@ -2528,7 +2582,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestPreExecute_EmptyExternalSource(
StorageConfig: &indexpb.StorageConfig{StorageType: "local"},
// ExternalSource is empty
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
err := task.PreExecute(ctx)
s.Error(err)
s.Contains(err.Error(), "external source is empty")
@@ -2545,7 +2599,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestPreExecute_InvalidExternalSpec(
Schema: &schemapb.CollectionSchema{},
StorageConfig: &indexpb.StorageConfig{StorageType: "local"},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
err := task.PreExecute(ctx)
s.Error(err)
s.Contains(err.Error(), "failed to parse external spec")
@@ -2562,7 +2616,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestPreExecute_UnsupportedFormat()
Schema: &schemapb.CollectionSchema{},
StorageConfig: &indexpb.StorageConfig{StorageType: "local"},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
err := task.PreExecute(ctx)
s.Error(err)
s.Contains(err.Error(), "unsupported format")
@@ -2751,7 +2805,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestFetchFragmentsFromExternalSourc
ExternalSpec: `{"format":"parquet"}`,
ExploreManifestPath: "", // empty
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
_, err := task.fetchFragmentsFromExternalSource(ctx)
@@ -2760,26 +2814,38 @@ func (s *RefreshExternalCollectionTaskSuite) TestFetchFragmentsFromExternalSourc
}
func (s *RefreshExternalCollectionTaskSuite) TestFetchFragmentsFromExternalSource_Success() {
paramtable.Init()
ctx := context.Background()
req := &datapb.RefreshExternalCollectionTaskRequest{
CollectionID: s.collectionID,
TaskID: s.taskID,
ExternalSource: "s3:///bucket/path",
ExternalSpec: `{"format":"parquet"}`,
ExploreManifestPath: "/manifests/explore.json",
FileIndexBegin: 0,
FileIndexEnd: 5,
StorageConfig: &indexpb.StorageConfig{StorageType: "local", BucketName: "/tmp"},
CollectionID: s.collectionID,
TaskID: s.taskID,
ExternalSource: "s3:///bucket/path",
ExternalSpec: `{"format":"parquet"}`,
ExploreManifestPath: "/manifests/explore.json",
FileIndexBegin: 0,
FileIndexEnd: 5,
StorageConfig: &indexpb.StorageConfig{StorageType: "local", BucketName: "/tmp"},
TargetRowsPerSegment: 12345,
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
task.columns = []string{"col1"}
mockFetch := mockey.Mock(packed.FetchFragmentsFromExternalSourceWithRange).
Return([]packed.Fragment{
{FragmentID: 0, FilePath: "f1.parquet", StartRow: 0, EndRow: 1000, RowCount: 1000},
}, nil).Build()
To(func(
ctx context.Context,
format string,
columns []string,
externalSource string,
storageConfig *indexpb.StorageConfig,
fileIndexBegin, fileIndexEnd int64,
exploreManifestPath string,
opts packed.ExternalFetchOptions,
) ([]packed.Fragment, error) {
s.Equal(req.GetTargetRowsPerSegment(), opts.RowLimit)
return []packed.Fragment{
{FragmentID: 0, FilePath: "f1.parquet", StartRow: 0, EndRow: 1000, RowCount: 1000},
}, nil
}).Build()
defer mockFetch.UnPatch()
frags, err := task.fetchFragmentsFromExternalSource(ctx)
@@ -2797,7 +2863,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Crea
ExternalSpec: `{"format":"parquet"}`,
StorageConfig: &indexpb.StorageConfig{StorageType: "local", BucketName: tmpDir, RootPath: tmpDir},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
task.columns = []string{"col1"}
task.preallocatedIDRange = &datapb.IDRange{Begin: 1, End: 100}
@@ -2828,7 +2894,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Cont
ExternalSpec: `{"format":"parquet"}`,
StorageConfig: &indexpb.StorageConfig{StorageType: "local", BucketName: tmpDir, RootPath: tmpDir},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
task.columns = []string{"col1"}
task.preallocatedIDRange = &datapb.IDRange{Begin: 1, End: 100}
@@ -2862,7 +2928,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_CtxC
ExternalSpec: `{"format":"parquet"}`,
StorageConfig: &indexpb.StorageConfig{StorageType: "local", BucketName: tmpDir, RootPath: tmpDir},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
task.columns = []string{"col1"}
task.preallocatedIDRange = &datapb.IDRange{Begin: 1, End: 100}
@@ -2997,7 +3063,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Samp
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -3051,7 +3117,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Samp
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -3101,7 +3167,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Samp
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -3158,7 +3224,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_PerS
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -3184,7 +3250,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_PerS
defer m2.UnPatch()
// 3 fragments each above targetRowsPerSegment → 3 segments.
targetRows := paramtable.Get().DataNodeCfg.ExternalCollectionTargetRowsPerSegment.GetAsInt64()
targetRows := req.GetTargetRowsPerSegment()
rowsPerFragment := targetRows * 2 // force one segment per fragment
fragments := []packed.Fragment{
{FragmentID: 1, RowCount: rowsPerFragment},
@@ -3227,7 +3293,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_PerS
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -3256,7 +3322,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_PerS
}).Build()
defer m2.UnPatch()
targetRows := paramtable.Get().DataNodeCfg.ExternalCollectionTargetRowsPerSegment.GetAsInt64()
targetRows := req.GetTargetRowsPerSegment()
rowsPerFragment := targetRows * 2
fragments := []packed.Fragment{
{FragmentID: 1, RowCount: rowsPerFragment},
@@ -3299,7 +3365,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_PerS
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -3314,7 +3380,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_PerS
Return(nil, fmt.Errorf("all samples failed")).Build()
defer m2.UnPatch()
targetRows := paramtable.Get().DataNodeCfg.ExternalCollectionTargetRowsPerSegment.GetAsInt64()
targetRows := req.GetTargetRowsPerSegment()
rowsPerFragment := targetRows * 2
fragments := []packed.Fragment{
{FragmentID: 1, RowCount: rowsPerFragment},
@@ -3353,7 +3419,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Zero
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -3413,7 +3479,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Insu
Schema: &schemapb.CollectionSchema{},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -3442,7 +3508,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Mani
Schema: &schemapb.CollectionSchema{},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -3477,7 +3543,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Cont
Schema: &schemapb.CollectionSchema{},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -3520,7 +3586,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Ensu
Schema: &schemapb.CollectionSchema{},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -3561,7 +3627,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Empt
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -3619,7 +3685,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Cont
ExternalSpec: `{"format":"parquet"}`,
Schema: &schemapb.CollectionSchema{},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -3654,7 +3720,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Cont
ExternalSpec: `{"format":"parquet"}`,
Schema: &schemapb.CollectionSchema{},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -3686,7 +3752,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Func
},
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -3812,7 +3878,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Adds
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
@@ -3865,7 +3931,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_Pass
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
// Parse the spec so the balance path has an in-memory struct to work with.
@@ -3933,7 +3999,7 @@ func (s *RefreshExternalCollectionTaskSuite) TestBalanceFragmentsToSegments_NilS
},
}
task := NewRefreshExternalCollectionTask(ctx, req)
task := s.newTask(ctx, req)
task.preallocatedIDRange = req.GetPreAllocatedSegmentIds()
task.nextAllocID = task.preallocatedIDRange.Begin
task.parsedSpec = &externalspec.ExternalSpec{Format: "parquet"}
+1
View File
@@ -1447,6 +1447,7 @@ message RefreshExternalCollectionTaskRequest {
int64 file_index_begin = 12; // File range start for this task
int64 file_index_end = 13; // File range end for this task
int64 partitionID = 14; // target partition for manifest path
int64 target_rows_per_segment = 15; // target row count for segment organization
}
message RefreshExternalCollectionTaskResponse {
File diff suppressed because it is too large Load Diff