@@ -255,24 +255,48 @@ pub(crate) trait OpenTsdbStorageReadExt: StorageRead {
255255 Ok ( max_series_id)
256256 }
257257
258- /// Load samples for a batch of series using a single sequential scan.
258+ /// Load samples for a batch of series using a narrowed sequential scan.
259259 ///
260- /// Instead of N individual `get()` calls, this scans the entire time series
261- /// key range for the bucket and filters to the requested series IDs. For
262- /// large candidate sets (e.g. 3796 series), a single sequential scan is
263- /// dramatically faster than thousands of random point reads because SST
264- /// blocks are read linearly with good cache locality.
265- #[ tracing:: instrument( level = "info" , skip( self , bucket, series_ids) , fields( bucket_start = bucket. start, wanted = series_ids. len( ) , scanned) ) ]
260+ /// Scans from the minimum to maximum series_id key, filtering to only
261+ /// the requested IDs via HashSet lookup. This is faster than N individual
262+ /// `get()` calls because sequential scan has much lower per-record overhead
263+ /// (~0.35ms vs ~6ms for point reads due to LSM tree traversal).
264+ ///
265+ /// Stops early once all requested series have been found.
266+ #[ tracing:: instrument( level = "info" , skip( self , bucket, series_ids) , fields( bucket_start = bucket. start, wanted = series_ids. len( ) , scanned, range_size) ) ]
266267 async fn get_time_series_batch (
267268 & self ,
268269 bucket : & TimeBucket ,
269270 series_ids : & [ SeriesId ] ,
270271 ) -> Result < HashMap < SeriesId , Vec < Sample > > > {
271272 let wanted: HashSet < SeriesId > = series_ids. iter ( ) . copied ( ) . collect ( ) ;
272- let range = TimeSeriesKey :: bucket_range ( bucket) ;
273+ let min_id = * series_ids. iter ( ) . min ( ) . unwrap ( ) ;
274+ let max_id = * series_ids. iter ( ) . max ( ) . unwrap ( ) ;
275+ let range_size = ( max_id - min_id + 1 ) as u64 ;
276+ tracing:: Span :: current ( ) . record ( "range_size" , range_size) ;
277+
278+ // Narrow scan from min to max+1 series_id key
279+ let start_key = TimeSeriesKey {
280+ time_bucket : bucket. start ,
281+ bucket_size : bucket. size ,
282+ series_id : min_id,
283+ }
284+ . encode ( ) ;
285+ let end_key = TimeSeriesKey {
286+ time_bucket : bucket. start ,
287+ bucket_size : bucket. size ,
288+ series_id : max_id. wrapping_add ( 1 ) ,
289+ }
290+ . encode ( ) ;
291+ let range = common:: BytesRange :: new (
292+ std:: ops:: Bound :: Included ( start_key) ,
293+ std:: ops:: Bound :: Excluded ( end_key) ,
294+ ) ;
295+
273296 let mut iter = self . scan_iter ( range) . await ?;
274297 let mut result = HashMap :: with_capacity ( series_ids. len ( ) ) ;
275298 let mut scanned = 0u64 ;
299+ let mut found = 0usize ;
276300 while let Some ( record) = iter. next ( ) . await ? {
277301 scanned += 1 ;
278302 let key = TimeSeriesKey :: decode ( record. key . as_ref ( ) ) ?;
@@ -282,33 +306,60 @@ pub(crate) trait OpenTsdbStorageReadExt: StorageRead {
282306 None => Vec :: new ( ) ,
283307 } ;
284308 result. insert ( key. series_id , samples) ;
309+ found += 1 ;
310+ if found == wanted. len ( ) {
311+ break ; // All requested series found
312+ }
285313 }
286314 }
287315 tracing:: Span :: current ( ) . record ( "scanned" , scanned) ;
288316 Ok ( result)
289317 }
290318
291- /// Load forward index entries for a batch of series using a single scan.
292- ///
293- /// Similar to `get_time_series_batch`, this uses a sequential scan instead
294- /// of individual gets when the candidate set is large.
295- #[ tracing:: instrument( level = "info" , skip( self , bucket, series_ids) , fields( bucket_start = bucket. start, wanted = series_ids. len( ) , scanned) ) ]
319+ /// Load forward index entries for a batch of series using a narrowed scan.
320+ #[ tracing:: instrument( level = "info" , skip( self , bucket, series_ids) , fields( bucket_start = bucket. start, wanted = series_ids. len( ) , scanned, range_size) ) ]
296321 async fn get_forward_index_batch (
297322 & self ,
298323 bucket : & TimeBucket ,
299324 series_ids : & [ SeriesId ] ,
300325 ) -> Result < ForwardIndex > {
301326 let wanted: HashSet < SeriesId > = series_ids. iter ( ) . copied ( ) . collect ( ) ;
302- let range = ForwardIndexKey :: bucket_range ( bucket) ;
327+ let min_id = * series_ids. iter ( ) . min ( ) . unwrap ( ) ;
328+ let max_id = * series_ids. iter ( ) . max ( ) . unwrap ( ) ;
329+ let range_size = ( max_id - min_id + 1 ) as u64 ;
330+ tracing:: Span :: current ( ) . record ( "range_size" , range_size) ;
331+
332+ let start_key = ForwardIndexKey {
333+ time_bucket : bucket. start ,
334+ bucket_size : bucket. size ,
335+ series_id : min_id,
336+ }
337+ . encode ( ) ;
338+ let end_key = ForwardIndexKey {
339+ time_bucket : bucket. start ,
340+ bucket_size : bucket. size ,
341+ series_id : max_id. wrapping_add ( 1 ) ,
342+ }
343+ . encode ( ) ;
344+ let range = common:: BytesRange :: new (
345+ std:: ops:: Bound :: Included ( start_key) ,
346+ std:: ops:: Bound :: Excluded ( end_key) ,
347+ ) ;
348+
303349 let mut iter = self . scan_iter ( range) . await ?;
304350 let result = ForwardIndex :: default ( ) ;
305351 let mut scanned = 0u64 ;
352+ let mut found = 0usize ;
306353 while let Some ( record) = iter. next ( ) . await ? {
307354 scanned += 1 ;
308355 let key = ForwardIndexKey :: decode ( record. key . as_ref ( ) ) ?;
309356 if wanted. contains ( & key. series_id ) {
310357 let value = ForwardIndexValue :: decode ( record. value . as_ref ( ) ) ?;
311358 result. series . insert ( key. series_id , value. into ( ) ) ;
359+ found += 1 ;
360+ if found == wanted. len ( ) {
361+ break ;
362+ }
312363 }
313364 }
314365 tracing:: Span :: current ( ) . record ( "scanned" , scanned) ;
0 commit comments