Parallel segment replay using scoped threads. Each segment is read in its own thread via mmap. Since segments are monotonically ordered by LSN, concatenating per-segment results in segment order produces a globally LSN-ordered result.
(
segments: &[crate::segment::SegmentMeta],
from_lsn: u64,
)
| 325 | /// monotonically ordered by LSN, concatenating per-segment results in |
| 326 | /// segment order produces a globally LSN-ordered result. |
| 327 | fn replay_segments_parallel( |
| 328 | segments: &[crate::segment::SegmentMeta], |
| 329 | from_lsn: u64, |
| 330 | ) -> Result<Vec<WalRecord>> { |
| 331 | // Collect per-segment results. Index corresponds to segment order. |
| 332 | let mut per_segment: Vec<Result<Vec<WalRecord>>> = Vec::with_capacity(segments.len()); |
| 333 | |
| 334 | std::thread::scope(|scope| { |
| 335 | let handles: Vec<_> = segments |
| 336 | .iter() |
| 337 | .map(|seg| { |
| 338 | scope.spawn(move || -> Result<Vec<WalRecord>> { |
| 339 | let mut reader = MmapWalReader::open(&seg.path)?; |
| 340 | let mut seg_records = Vec::new(); |
| 341 | while let Some(record) = reader.next_record()? { |
| 342 | if record.header.lsn >= from_lsn { |
| 343 | seg_records.push(record); |
| 344 | } |
| 345 | } |
| 346 | reader.release_pages(); |
| 347 | Ok(seg_records) |
| 348 | }) |
| 349 | }) |
| 350 | .collect(); |
| 351 | |
| 352 | for handle in handles { |
| 353 | per_segment.push(handle.join().unwrap_or_else(|_| { |
| 354 | Err(WalError::Io(std::io::Error::other( |
| 355 | "segment replay thread panicked", |
| 356 | ))) |
| 357 | })); |
| 358 | } |
| 359 | }); |
| 360 | |
| 361 | // Merge in segment order (preserves LSN ordering). |
| 362 | let total_estimate: usize = per_segment |
| 363 | .iter() |
| 364 | .map(|r| r.as_ref().map(|v| v.len()).unwrap_or(0)) |
| 365 | .sum(); |
| 366 | let mut records = Vec::with_capacity(total_estimate); |
| 367 | for seg_result in per_segment { |
| 368 | records.extend(seg_result?); |
| 369 | } |
| 370 | |
| 371 | Ok(records) |
| 372 | } |
| 373 | |
| 374 | /// Paginated mmap replay: reads at most `max_records` from `from_lsn`. |
| 375 | /// |
no test coverage detected