docs/design-docs/design_docs/20260129-add-function-field-design.md
Commit: 513c92d7f2 feat: support add function field (#44444) Author: MrPresent-Han Date: September 2025 Scope: 156 files, +10,129/-3,475 lines
Function fields enable users to dynamically add computed/derived fields to existing collections without requiring data re-ingestion. The primary use case is adding BM25 (sparse vector) fields to collections that were originally created with only dense vector fields, enabling hybrid search capabilities post-creation.
┌─────────────────────────────────────────────────────────────────────────────┐
│ AlterCollectionSchema Request │
└─────────────────────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────────────────┐
│ PROXY │
│ • Validate request (field types, names, schema version consistency) │
│ • Check all segments have aligned schema versions │
│ • Create alterCollectionSchemaTask and enqueue to DDL queue │
│ • Optionally create indexes for new fields │
└─────────────────────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────────────────┐
│ ROOTCOORD │
│ • Validate function schema (input/output fields, uniqueness) │
│ • Assign field IDs and function IDs │
│ • Increment schema version │
│ • Broadcast AlterCollectionMessage to WAL + all virtual channels │
└─────────────────────────────────────────────────────────────────────────────┘
│
┌───────────────────┼───────────────────┐
▼ ▼ ▼
┌──────────────────────────┐ ┌──────────────────────┐ ┌──────────────────────┐
│ STREAMING/WAL │ │ DATACOORD │ │ QUERYNODE │
│ • Flush existing segments│ │ • Track segment │ │ • SyncSchema to │
│ • Log schema change │ │ schema versions │ │ segments │
│ • Update in-memory schema│ │ • Trigger backfill │ │ • Update IDF oracle │
│ • Validate insert schema │ │ compaction if │ │ • Rebuild function │
│ versions │ │ enabled │ │ runners │
└──────────────────────────┘ └──────────────────────┘ └──────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────────────────┐
│ DATANODE (Backfill) │
│ • Execute backfill compaction for segments with outdated schema │
│ • Compute function outputs (e.g., BM25 sparse vectors) │
│ • Write new binlogs with function field data │
│ • Update BM25 statistics │
└─────────────────────────────────────────────────────────────────────────────┘
| Component | Responsibility |
|---|---|
| Proxy | API gateway, request validation, schema consistency checks, task orchestration |
| RootCoord | Schema management, ID assignment, metadata persistence, broadcast coordination |
| Streaming/WAL | Durable schema change logging, write consistency, version mismatch detection |
| DataCoord | Segment metadata tracking, backfill compaction policy, schema version monitoring |
| DataNode | Backfill compaction execution, function output computation, binlog writing |
| QueryNode | Schema synchronization, function runner management, IDF oracle updates |
| Segcore (C++) | Low-level schema sync, field accessibility checks, data storage |
rpc AlterCollectionSchema(AlterCollectionSchemaRequest) returns (AlterCollectionSchemaResponse) {}
Request Structure:
db_name: Database namecollection_name: Collection nameaction: Schema alteration action (currently only ADD supported)add_request: Contains field infos and function schemas to adddo_physical_backfill: Whether to backfill existing dataConstraints:
rpc UpdateIndex(UpdateIndexRequest) returns (common.Status) {}
message UpdateIndexRequest {
common.MsgBase base = 1;
int64 collectionID = 2;
oneof Action {
AddIndex add_index_request = 3;
DropIndex drop_index_request = 4;
}
}
Multiple message types now include schema_version for tracking:
// data_coord.proto
message AllocSegmentRequest {
int32 schema_version = 7;
}
message SegmentInfo {
int32 schema_version = 33;
}
// messages.proto
message InsertMessageHeader {
int32 schema_version = 3;
}
message CreateSegmentMessageHeader {
int32 schema_version = 8;
}
// data_coord.proto
enum CompactionType {
BackfillCompaction = 12;
}
message CompactionPlan {
repeated schema.FunctionSchema functions = 30;
}
message CompactionTask {
repeated schema.FunctionSchema diff_functions = 29;
}
// index_coord.proto
message FieldIndex {
int32 min_schema_version = 4;
}
// streaming.proto
enum StreamingCode {
STREAMING_CODE_SCHEMA_VERSION_MISMATCH = 14;
}
File: internal/proxy/impl.go, internal/proxy/task.go
1. Health Check
│
▼
2. DescribeCollection (get current schema)
│
▼
3. Schema Version Consistency Check
• GetCollectionStatistics
• Verify SchemaVersionConsistencyProportion == 100%
│
▼
4. Create alterCollectionSchemaTask
│
▼
5. Enqueue to DDL Queue
│
▼
6. PreExecute: Validate fields and function schema
• Check max field count
• Check no duplicate names
• Validate data types
• Check not system field
│
▼
7. Execute: Call RootCoord.AlterCollectionSchema
│
▼
8. Post-Execute: Create indexes if configured
Before allowing schema changes, the proxy validates that all segments have aligned schema versions:
stats, _ := s.GetCollectionStatistics(ctx, collectionID)
proportion := stats[common.SchemaVersionConsistencyProportionKey]
if proportion != "100" {
return errors.New("segments have inconsistent schema versions")
}
File: internal/rootcoord/ddl_callbacks_alter_collection_schema.go
1. Acquire broadcast lock on collection
│
▼
2. Retrieve current collection metadata
│
▼
3. Validate request:
• Exactly one function schema
• Valid field schemas
• No duplicate field names
• Function name doesn't exist
│
▼
4. Assign IDs:
• Field IDs from nextFieldID(coll)
• Function ID from nextFunctionID(coll)
• Resolve field names → field IDs for function I/O
│
▼
5. Construct new schema:
• Copy existing fields and functions
• Append new fields and function
• Increment schema version
• Set DoPhysicalBackfill flag
│
▼
6. Broadcast AlterCollectionMessage:
• Send to control channel
• Send to all virtual channels
// Field ID assignment
func nextFieldID(coll *model.Collection) int64 {
maxFieldID := findMaxFieldID(coll.Fields, coll.StructArraySubFields)
return maxFieldID + 1
}
// Function ID assignment
func nextFunctionID(coll *model.Collection) int64 {
maxFunctionID := common.StartOfUserFunctionID
for _, fn := range coll.Functions {
if fn.ID > maxFunctionID {
maxFunctionID = fn.ID
}
}
return maxFunctionID + 1
}
Files: internal/streamingnode/server/wal/interceptors/shard/
The shard interceptor validates schema versions at write time:
func handleInsertMessage(ctx context.Context, msg InsertMessage) error {
schemaVersion := msg.Header.GetSchemaVersion()
correctVersion, err := shardManager.CheckIfCollectionSchemaVersionMatch(
msg.Header.GetCollectionId(),
schemaVersion,
)
if err != nil {
return status.NewSchemaVersionMismatch(
"schema version mismatch, input: %d, collection: %d",
schemaVersion, correctVersion,
)
}
// Process insert...
}
Schema changes follow strict ordering to maintain consistency:
func handleAlterCollection(ctx context.Context, msg AlterCollectionMessage) error {
// 1. FLUSH existing segments FIRST (creates checkpoint)
if messageutil.IsSchemaChange(header) {
segmentIDs, _ := flushSegments(ctx, collectionID)
header.FlushedSegmentIds = segmentIDs
}
// 2. APPEND to WAL (durable record)
msgID, err := appendOp(ctx, msg)
if err != nil {
return err
}
// 3. UPDATE in-memory state LAST (after WAL success)
alterCollectionMsg := message.AsImmutableAlterCollectionMessageV2(msg)
if err := shardManager.AlterCollection(alterCollectionMsg); err != nil {
panic("failed to alter collection after WAL append")
}
}
Why panic() is Used (Critical Design Decision):
The panic at step 3 is intentional and represents an unrecoverable state where:
WAL-Memory Inconsistency: The schema change has been durably written to WAL but failed to apply to in-memory state. This creates a dangerous inconsistency where:
Why Alternatives Don't Work:
Recovery Process After Panic:
AlterCollectionMessage updates in-memory state correctlyConsistency Guarantees:
This ordering ensures:
Files: internal/datacoord/meta.go, internal/datacoord/compaction_policy_backfill.go
Each segment tracks its schema version:
type SegmentInfo struct {
// ... other fields
SchemaVersion int32
}
The DataCoord calculates schema consistency metrics:
func GetCollectionStatistics(ctx context.Context, req Request) Response {
collectionSchemaVersion := collection.Schema.GetVersion()
segments := meta.SelectSegments(ctx, WithCollection(req.CollectionID))
consistentCount := 0
for _, segment := range segments {
if segment.GetSchemaVersion() == collectionSchemaVersion {
consistentCount++
}
}
proportion := float64(consistentCount) / float64(len(segments)) * 100.0
return Response{
Stats: map[string]string{
SchemaVersionConsistencyProportionKey: fmt.Sprintf("%.2f", proportion),
},
}
}
Trigger Conditions:
DoPhysicalBackfill = trueSchemaVersion < Collection.SchemaVersionPolicy Flow:
func (p *backfillCompactionPolicy) Trigger(ctx context.Context) ([]CompactionView, error) {
for _, collection := range collections {
segments := getEligibleSegments(collection.ID)
for _, segment := range segments {
if segment.SchemaVersion < collection.SchemaVersion {
if collection.DoPhysicalBackfill {
// Get schema diff to identify new functions
oldSchema := getSchemaByVersion(segment.SchemaVersion)
funcDiff := util.SchemaDiff(oldSchema, collection.Schema)
// Create backfill compaction view
views = append(views, BackfillSegmentsView{
segmentID: segment.ID,
funcDiff: funcDiff,
})
} else {
// Just update metadata (no physical backfill)
segment.SchemaVersion = collection.SchemaVersion
}
}
}
}
return views, nil
}
File: internal/datanode/compactor/backfill_compactor.go
┌─────────────────────────────────────────────────────────────────┐
│ Backfill Compaction Pipeline │
├─────────────────────────────────────────────────────────────────┤
│ │
│ 1. Pre-Validation │
│ • Exactly one segment in plan │
│ • Field binlogs present │
│ • Exactly one backfill function │
│ • FunctionRunner validates successfully │
│ │
│ 2. Read Input Data │
│ • Read input field binlogs (e.g., varchar for BM25) │
│ • Decompress and parse via BinlogRecordReader │
│ • Build input data array │
│ │
│ 3. Execute Function │
│ • Run FunctionRunner.BatchRun() on input data │
│ • For BM25: Compute sparse float vectors │
│ • Build InsertData with function outputs │
│ │
│ 4. Write Output │
│ • Create PackedWriter for new field binlogs │
│ • Allocate new log IDs │
│ • Write records to object storage │
│ │
│ 5. Update Statistics (BM25) │
│ • Serialize BM25 stats (term frequencies) │
│ • Write to dedicated BM25 stats log files │
│ │
│ 6. Merge Logs │
│ • Combine new function field binlogs with original binlogs │
│ • Create FieldBinlog entries with sizes │
│ │
│ 7. Return Result │
│ • CompactionPlanResult with merged logs │
│ • Segment ID, row count, BM25 logs │
│ │
└─────────────────────────────────────────────────────────────────┘
type backfillMetrics struct {
getInputDataDuration time.Duration
executeBM25Duration time.Duration
writeRecordDuration time.Duration
updateStatsDuration time.Duration
}
Files: internal/querynodev2/delegator/delegator.go, internal/querynodev2/pipeline/embedding_node.go
func (sd *shardDelegator) UpdateSchema(ctx context.Context, schema *schemapb.CollectionSchema) error {
// Update collection manager
sd.collection.UpdateSchema(schema)
// Update BM25 function runners
sd.updateBM25Functions(schema, ctx)
// Propagate to all segments
return sd.propagateSchemaToSegments(ctx, schema)
}
func (sd *shardDelegator) updateBM25Functions(schema *schemapb.CollectionSchema, ctx context.Context) {
// Get current BM25 output field IDs
currentOutputFields := getCurrentBM25OutputFields(sd.schema)
// Get new BM25 output field IDs
newOutputFields := getBM25OutputFields(schema)
// Find only NEW functions (not in current set)
for fieldID := range newOutputFields {
if _, exists := currentOutputFields[fieldID]; !exists {
// Create function runner for new BM25 function
runner := createFunctionRunner(schema, fieldID)
sd.functionRunners[fieldID] = runner
sd.analyzerRunners[inputFieldID] = runner
sd.isBM25Field[fieldID] = true
}
}
// Update or create IDF Oracle
if sd.idfOracle == nil {
sd.idfOracle = NewIDFOracle(schema.Functions)
} else {
sd.idfOracle.UpdateCurrent(schema.Functions)
}
}
The embedding node dynamically adapts to schema changes:
type embeddingNode struct {
curSchema *schemapb.CollectionSchema
functionRunners map[int64]function.FunctionRunner // keyed by function ID
}
func (en *embeddingNode) Operate(msgs []flowgraph.Msg) []flowgraph.Msg {
for _, msg := range msgs {
insertMsg := msg.(*insertNodeMsg)
// Check for schema update
if insertMsg.schema != nil && insertMsg.schema != en.curSchema {
en.curSchema = insertMsg.schema
en.setupFunctionRunners() // Rebuild runners for new schema
}
// Process with current function runners
if len(en.functionRunners) > 0 {
en.processWithFunctions(insertMsg)
}
}
}
Files: internal/core/src/common/Schema.h, internal/core/src/segcore/SegmentInterface.h
New SyncSchema() operation allows runtime schema updates:
class SegmentInternalInterface {
public:
void SyncSchema(SchemaPtr new_schema) {
std::unique_lock<std::shared_mutex> lock(sch_mutex_);
if (new_schema->get_schema_version() > schema_->get_schema_version()) {
schema_ = new_schema;
}
}
protected:
SchemaPtr schema_;
mutable std::shared_mutex sch_mutex_; // Thread-safe schema access
};
New method to check if a field is accessible (either has data or index):
bool FieldAccessible(FieldId field_id) const {
return HasFieldData(field_id) || HasIndex(field_id);
}
Vector search operations now gracefully handle missing function fields:
std::unique_ptr<SearchResult> AsyncSearch(SearchInfo& search_info) {
FieldId target_field = search_info.GetFieldId();
// Check if function field is accessible
if (!segment->FieldAccessible(target_field)) {
// Return empty result instead of crashing
return std::make_unique<SearchResult>(
make_empty_search_result(search_info)
);
}
// Proceed with normal search
return DoSearch(search_info);
}
File: internal/util/schema_util.go
type FieldDiff struct {
Added []*schemapb.FieldSchema // Fields in new but not in old
}
type FuncDiff struct {
Added []*schemapb.FunctionSchema // Functions in new but not in old
}
func SchemaDiff(oldSchema, newSchema *schemapb.CollectionSchema) (*FieldDiff, *FuncDiff, error) {
if oldSchema == nil || newSchema == nil {
return nil, nil, errors.New("schema cannot be nil")
}
fieldDiff := compareFields(oldSchema.Fields, newSchema.Fields)
funcDiff := compareFunctions(oldSchema.Functions, newSchema.Functions)
return fieldDiff, funcDiff, nil
}
func compareFunctions(oldFuncs, newFuncs []*schemapb.FunctionSchema) *FuncDiff {
// Build map of old function IDs for O(1) lookup
oldMap := make(map[int64]bool)
for _, fn := range oldFuncs {
if fn != nil {
oldMap[fn.Id] = true
}
}
// Find functions in new but not in old
var added []*schemapb.FunctionSchema
for _, fn := range newFuncs {
if fn != nil && !oldMap[fn.Id] {
added = append(added, fn)
}
}
return &FuncDiff{Added: added}
}
File: internal/querynodev2/delegator/idf_oracle.go
New method to handle function field additions:
func (oracle *IDFOracle) UpdateCurrent(functions []*schemapb.FunctionSchema) {
oracle.mu.Lock()
defer oracle.mu.Unlock()
for _, fn := range functions {
if fn.Type == schemapb.FunctionType_BM25 {
outputFieldID := fn.OutputFieldIds[0]
// Initialize stats for new BM25 fields
if _, exists := oracle.currentStats[outputFieldID]; !exists {
oracle.currentStats[outputFieldID] = NewBM25Stats()
}
}
}
}
Enhanced merging for backfilled segments:
func (seg *segmentStats) MergeStats(newStats bm25Stats) bool {
seg.mu.Lock()
defer seg.mu.Unlock()
// Load from disk if needed
if seg.stats == nil && seg.statsPath != "" {
seg.stats = loadStatsFromLocalNoLock(seg.statsPath)
}
// Merge stats
for fieldID, newFieldStats := range newStats {
if oldStats, exists := seg.stats[fieldID]; exists {
oldStats.Merge(newFieldStats)
} else {
seg.stats[fieldID] = newFieldStats.Clone()
}
}
return seg.activated // Return whether to update current stats
}
File: internal/datacoord/index_service.go
Indexes now track the minimum schema version required:
func (s *Server) CreateIndex(ctx context.Context, req *indexpb.CreateIndexRequest) error {
// Get latest schema
schema, _ := s.broker.DescribeCollectionInternal(ctx, collectionID, typeutil.MaxTimestamp)
index := &model.Index{
// ... other fields
MinSchemaVersion: schema.GetVersion(), // NEW: Track schema version
}
// Broadcast to all channels including control channel
channels := append([]string{streaming.WAL().ControlChannel()}, vchannels...)
return s.saveAndBroadcastIndex(ctx, index, channels)
}
New streaming error code for version conflicts:
const STREAMING_CODE_SCHEMA_VERSION_MISMATCH = 14
func (e *StreamingError) IsSchemaVersionMismatch() bool {
return e.Code == streamingpb.StreamingCode_STREAMING_CODE_SCHEMA_VERSION_MISMATCH
}
func (e *StreamingError) IsUnrecoverable() bool {
return e.Code == STREAMING_CODE_UNRECOVERABLE ||
e.IsReplicateViolation() ||
e.IsTxnUnavailable() ||
e.IsSchemaVersionMismatch() // Schema mismatches are unrecoverable
}
Segments without function field data return empty results rather than failing:
// In VectorSearchNode
if (!segment->FieldAccessible(target_vector_field_id)) {
return make_empty_search_result(num_queries, topK);
}
// component_param.go
type BackfillConfig struct {
// Whether backfill compaction is enabled
Enabled bool
// Maximum concurrent backfill tasks
MaxConcurrentTasks int
// Backfill batch size
BatchSize int
}
Client Proxy RootCoord Streaming DataCoord QueryNode
│ │ │ │ │ │
│ AlterCollectionSchema │ │ │ │
├──────────────────►│ │ │ │ │
│ │ DescribeCollection │ │ │
│ ├────────────────►│ │ │ │
│ │◄────────────────┤ │ │ │
│ │ │ │ │ │
│ │ GetCollectionStatistics │ │ │
│ ├─────────────────────────────────────────────────►│ │
│ │◄─────────────────────────────────────────────────┤ │
│ │ (check schema version consistency = 100%) │ │
│ │ │ │ │ │
│ │ AlterCollectionSchema │ │ │
│ ├────────────────►│ │ │ │
│ │ │ │ │ │
│ │ │ Broadcast AlterCollectionMessage │ │
│ │ ├───────────────►│ │ │
│ │ │ │ │ │
│ │ │ (handleAlterCollection) │ │
│ │ │ • Flush existing segments │ │
│ │ │ • Append to WAL │ │
│ │ │ • Update in-memory state │ │
│ │ │ │ │ │
│ │ │ │ Forward to │ │
│ │ │ ├────────────────►│ │
│ │ │ │ DataCoord │ │
│ │ │ │ │ │
│ │ │ │ UpdateSchema │ │
│ │ │ ├────────────────────────────────►│
│ │ │ │ │ │
│ │◄────────────────┤ │ │ │
│◄──────────────────┤ │ │ │ │
│ │ │ │ │ │
Key Points:
AlterCollectionMessage to the Streaming nodehandleAlterCollection function internally performs three steps in strict order:
DataCoord DataNode ObjectStorage
│ │ │
│ (Backfill policy detects │
│ segment with old schema) │
│ │ │
│ SubmitBackfillCompaction│ │
├────────────────────────►│ │
│ │ │
│ │ Read input field binlogs │
│ ├─────────────────────────►│
│ │◄─────────────────────────┤
│ │ │
│ │ Execute BM25 function │
│ │ (compute sparse vectors) │
│ │ │
│ │ Write output binlogs │
│ ├─────────────────────────►│
│ │◄─────────────────────────┤
│ │ │
│ │ Write BM25 stats │
│ ├─────────────────────────►│
│ │◄─────────────────────────┤
│ │ │
│ CompactionPlanResult │ │
│◄────────────────────────┤ │
│ │ │
│ Update segment metadata │ │
│ (schemaVersion = new) │ │
│ │ │
Decision: Use monotonically increasing integer version numbers.
Rationale:
<, >, ==)Decision: Support both modes via DoPhysicalBackfill flag.
Physical Backfill (DoPhysicalBackfill = true):
Logical Backfill (DoPhysicalBackfill = false):
Decision: Limit to one function field addition per request.
Rationale:
Decision: Flush segments before schema changes, log to WAL before updating in-memory state.
Rationale:
Decision: Return empty results for inaccessible function fields instead of failing.
Rationale:
| Component | Test File | Coverage |
|---|---|---|
| Schema Util | internal/util/schema_util_test.go | Field diff, function diff, nil handling |
| Backfill Policy | internal/datacoord/compaction_policy_backfill_test.go | Trigger conditions, segment selection |
| Backfill Task | internal/datacoord/compaction_task_backfill_test.go | State machine, progress tracking |
| Backfill Compactor | internal/datanode/compactor/backfill_compactor_test.go | Execution pipeline, error handling |
Support adding multiple function fields in a single request for efficiency.
Support modifying function parameters without full re-computation.
Support removing function fields with proper cleanup of binlogs and indexes.
Support pausing and resuming backfill operations for large collections.
Extend beyond BM25 to support user-defined function types.