Skip to content

Commit 36ac38b

Browse files
authored
[ENH] Skip record load when only id is requested (#6540)
## Description of changes _Summarize the changes made by this PR._ - Improvements & Bug fixes - Update `Select` and `Projection` operator so that if only user id is requested as output, skip loading the record data - New functionality - Re-introduce `get_user_id_for_offset_id` and `load_id_to_user_id` in record segment reader - Introduce `get_user_id` on materialized log ## Test plan _How are these changes tested?_ - [ ] Tests pass locally with `pytest` for python, `yarn test` for js, `cargo test` for rust ## Migration plan _Are there any migrations, or any forwards/backwards compatibility changes needed in order to make sure this change deploys reliably?_ ## Observability plan _What is the plan to instrument and monitor this change?_ ## Documentation Changes _Are all docstrings for user-facing APIs updated if required? Do we need to make documentation changes in the [docs section](https://github.com/chroma-core/chroma/tree/main/docs/docs.trychroma.com)?_
1 parent fb8c31f commit 36ac38b

4 files changed

Lines changed: 377 additions & 124 deletions

File tree

rust/segment/src/blockfile_record.rs

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -940,6 +940,28 @@ impl RecordSegmentReader<'_> {
940940
.await
941941
}
942942

943+
/// Get the user id for a given offset id using the lightweight id_to_user_id blockfile.
944+
/// This avoids loading the full DataRecord (embedding, metadata, document).
945+
/// Returns an error if the offset id is not found.
946+
pub async fn get_user_id_for_offset_id(
947+
&self,
948+
offset_id: u32,
949+
) -> Result<&str, Box<dyn ChromaError>> {
950+
self.id_to_user_id.get("", offset_id).await?.ok_or_else(|| {
951+
Box::new(RecordSegmentReaderCreationError::DataRecordNotFound(
952+
offset_id,
953+
)) as Box<dyn ChromaError>
954+
})
955+
}
956+
957+
/// Bulk prefetch for the id_to_user_id blockfile.
958+
/// This is the lightweight alternative to load_id_to_data when only user IDs are needed.
959+
pub async fn load_id_to_user_id(&self, keys: impl Iterator<Item = u32>) {
960+
self.id_to_user_id
961+
.load_data_for_keys(keys.map(|k| ("".to_string(), k)))
962+
.await
963+
}
964+
943965
pub async fn get_total_logical_size_bytes(&self) -> Result<u64, Box<dyn ChromaError>> {
944966
self.id_to_data
945967
.get_range_stream(""..="", ..)

rust/segment/src/types.rs

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -106,6 +106,10 @@ pub enum LogMaterializerError {
106106
EmbeddingMaterialization,
107107
#[error("Error reading record segment {0}")]
108108
RecordSegment(#[from] Box<dyn ChromaError>),
109+
#[error("Log index {0} out of bounds when resolving user ID")]
110+
LogIndexOutOfBounds(usize),
111+
#[error("Record segment reader required but not available")]
112+
RecordSegmentReaderRequired,
109113
}
110114

111115
impl ChromaError for LogMaterializerError {
@@ -114,6 +118,8 @@ impl ChromaError for LogMaterializerError {
114118
LogMaterializerError::MetadataMaterialization(e) => e.code(),
115119
LogMaterializerError::EmbeddingMaterialization => ErrorCodes::Internal,
116120
LogMaterializerError::RecordSegment(e) => e.code(),
121+
LogMaterializerError::LogIndexOutOfBounds(_) => ErrorCodes::Internal,
122+
LogMaterializerError::RecordSegmentReaderRequired => ErrorCodes::Internal,
117123
}
118124
}
119125
}
@@ -262,6 +268,32 @@ impl<'log_data> BorrowedMaterializedLogRecord<'log_data> {
262268
}
263269
}
264270

271+
/// Get the user id for this record without hydrating the full data record.
272+
/// When the user id is available in the log, it is returned directly (no I/O).
273+
/// Otherwise, it falls back to the lightweight id_to_user_id blockfile lookup.
274+
pub async fn get_user_id(
275+
&self,
276+
record_segment_reader: Option<&RecordSegmentReader<'_>>,
277+
) -> Result<String, LogMaterializerError> {
278+
if let Some(id) = self.materialized_log_record.user_id_at_log_index {
279+
return Ok(self
280+
.logs
281+
.get(id)
282+
.ok_or(LogMaterializerError::LogIndexOutOfBounds(id))?
283+
.record
284+
.id
285+
.clone());
286+
}
287+
288+
match record_segment_reader {
289+
Some(reader) => Ok(reader
290+
.get_user_id_for_offset_id(self.materialized_log_record.offset_id)
291+
.await?
292+
.to_string()),
293+
None => Err(LogMaterializerError::RecordSegmentReaderRequired),
294+
}
295+
}
296+
265297
/// Reads any record segment data that this log record may reference and returns a hydrated version of this record.
266298
/// The record segment reader passed here **must be over the same set of blockfiles** as the reader that was originally passed to `materialize_logs()`. If the two readers are different, the behavior is undefined.
267299
pub async fn hydrate<'segment_data>(

rust/worker/src/execution/operators/projection.rs

Lines changed: 92 additions & 48 deletions
Original file line numberDiff line numberDiff line change
@@ -73,9 +73,12 @@ impl Operator<ProjectionInput, ProjectionOutput> for Projection {
7373
return Ok(ProjectionOutput { records: vec![] });
7474
}
7575

76+
let needs_data = self.document || self.embedding || self.metadata;
77+
7678
tracing::trace!(
77-
"Running projection on {} offset ids",
78-
input.offset_ids.len()
79+
"Running projection on {} offset ids (needs_data={})",
80+
input.offset_ids.len(),
81+
needs_data,
7982
);
8083
let record_segment_reader = match Box::pin(RecordSegmentReader::from_segment(
8184
&input.record_segment,
@@ -90,12 +93,28 @@ impl Operator<ProjectionInput, ProjectionOutput> for Projection {
9093
Err(e) => Err(*e),
9194
}?;
9295

96+
let offset_id_set: HashSet<_> = HashSet::from_iter(input.offset_ids.iter().cloned());
97+
98+
// Prefetch: when needs_data, load full data records (which contain the user ID);
99+
// otherwise, load only the lightweight id_to_user_id mapping.
100+
if let Some(reader) = &record_segment_reader {
101+
if needs_data {
102+
reader
103+
.load_id_to_data(offset_id_set.iter().cloned())
104+
.instrument(tracing::trace_span!(parent: Span::current(), "Load ID to data", num_ids = offset_id_set.len()))
105+
.await;
106+
} else {
107+
reader
108+
.load_id_to_user_id(offset_id_set.iter().cloned())
109+
.instrument(tracing::trace_span!(parent: Span::current(), "Load ID to user ID", num_ids = offset_id_set.len()))
110+
.await;
111+
}
112+
}
113+
93114
let materialized_logs = materialize_logs(&record_segment_reader, input.logs.clone(), None)
94115
.instrument(tracing::trace_span!(parent: Span::current(), "Materialize logs"))
95116
.await?;
96117

97-
let offset_id_set: HashSet<_> = HashSet::from_iter(input.offset_ids.iter().cloned());
98-
99118
// Create a hash map that maps an offset id to the corresponding log
100119
// It contains all records from the logs that should be present in the final result
101120
let offset_id_to_log_record: HashMap<_, _> = materialized_logs
@@ -113,51 +132,76 @@ impl Operator<ProjectionInput, ProjectionOutput> for Projection {
113132
.iter()
114133
.map(|offset_id| {
115134
async {
116-
let record = match offset_id_to_log_record.get(offset_id) {
117-
// The offset id is in the log
118-
Some(log) => {
119-
let log = log
120-
.hydrate(record_segment_reader.as_ref())
121-
.await
122-
.map_err(ProjectionError::LogMaterializer)?;
123-
124-
ProjectionRecord {
125-
id: log.get_user_id().to_string(),
126-
document: log
127-
.merged_document_ref()
128-
.filter(|_| self.document)
129-
.map(str::to_string),
130-
embedding: self
131-
.embedding
132-
.then_some(log.merged_embeddings_ref().to_vec()),
133-
metadata: self
134-
.metadata
135-
.then_some(log.merged_metadata())
136-
.filter(|metadata| !metadata.is_empty()),
137-
}
138-
}
139-
// The offset id is in the record segment
140-
None => {
141-
if let Some(reader) = &record_segment_reader {
142-
let record =
143-
reader.get_data_for_offset_id(*offset_id).await?.ok_or(
144-
ProjectionError::RecordSegmentPhantomRecord(*offset_id),
145-
)?;
146-
ProjectionRecord {
147-
id: record.id.to_string(),
148-
document: record
149-
.document
150-
.filter(|_| self.document)
151-
.map(str::to_string),
152-
embedding: self.embedding.then_some(record.embedding.to_vec()),
153-
metadata: record.metadata.filter(|_| self.metadata),
135+
if needs_data {
136+
// Full hydration path: get ID and data from hydrated record
137+
let (id, document, embedding, metadata) =
138+
match offset_id_to_log_record.get(offset_id) {
139+
Some(log) => {
140+
let log = log
141+
.hydrate(record_segment_reader.as_ref())
142+
.await
143+
.map_err(ProjectionError::LogMaterializer)?;
144+
(
145+
log.get_user_id().to_string(),
146+
log.merged_document_ref()
147+
.filter(|_| self.document)
148+
.map(str::to_string),
149+
self.embedding
150+
.then_some(log.merged_embeddings_ref().to_vec()),
151+
self.metadata
152+
.then_some(log.merged_metadata())
153+
.filter(|metadata| !metadata.is_empty()),
154+
)
155+
}
156+
None => {
157+
let reader = record_segment_reader
158+
.as_ref()
159+
.ok_or(ProjectionError::RecordSegmentUninitialized)?;
160+
let record =
161+
reader.get_data_for_offset_id(*offset_id).await?.ok_or(
162+
ProjectionError::RecordSegmentPhantomRecord(*offset_id),
163+
)?;
164+
(
165+
record.id.to_string(),
166+
record
167+
.document
168+
.filter(|_| self.document)
169+
.map(str::to_string),
170+
self.embedding.then_some(record.embedding.to_vec()),
171+
record.metadata.filter(|_| self.metadata),
172+
)
154173
}
155-
} else {
156-
return Err(ProjectionError::RecordSegmentUninitialized);
157-
}
158-
}
159-
};
160-
Ok::<_, ProjectionError>(record)
174+
};
175+
176+
Ok(ProjectionRecord {
177+
id,
178+
document,
179+
embedding,
180+
metadata,
181+
})
182+
} else {
183+
// Lightweight path: resolve user ID only via id_to_user_id blockfile
184+
let id = match offset_id_to_log_record.get(offset_id) {
185+
Some(log) => log
186+
.get_user_id(record_segment_reader.as_ref())
187+
.await
188+
.map_err(ProjectionError::LogMaterializer)?,
189+
None => match &record_segment_reader {
190+
Some(reader) => reader
191+
.get_user_id_for_offset_id(*offset_id)
192+
.await?
193+
.to_string(),
194+
None => return Err(ProjectionError::RecordSegmentUninitialized),
195+
},
196+
};
197+
198+
Ok(ProjectionRecord {
199+
id,
200+
document: None,
201+
embedding: None,
202+
metadata: None,
203+
})
204+
}
161205
}
162206
.instrument(current_span.clone())
163207
})

0 commit comments

Comments
 (0)