Skip to content
Open
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
181 changes: 154 additions & 27 deletions qdp/qdp-core/src/pipeline_runner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
// Throughput/latency pipeline using QdpEngine and encode_batch. Full loop runs in Rust;
// Python bindings release GIL during the run.

use std::collections::VecDeque;
use std::f64::consts::PI;
use std::path::Path;
use std::time::Instant;
Expand Down Expand Up @@ -318,8 +319,9 @@ impl<T: FloatElem + ToBatchData> BatchProducer for InMemoryProducer<T> {

pub struct StreamingProducer<T: FloatElem = f64> {
pub reader: ParquetStreamingReader<T>,
pub buffer: Vec<T>,
pub buffer_cursor: usize,
/// Ring buffer of not-yet-consumed elements. Discarding a consumed prefix is O(1): a
/// `VecDeque` advances its head instead of shifting the retained tail like `Vec::drain`.
pub buffer: VecDeque<T>,
pub read_chunk_scratch: Vec<T>,
pub sample_size: usize,
pub batch_size: usize,
Expand All @@ -334,40 +336,46 @@ impl<T: FloatElem + ToBatchData> BatchProducer for StreamingProducer<T> {
return Ok(None);
}
let required = self.batch_size * self.sample_size;
while (self.buffer.len() - self.buffer_cursor) < required {
while self.buffer.len() < required {
let written = self.reader.read_chunk(&mut self.read_chunk_scratch)?;
if written == 0 {
break;
}
self.buffer
.extend_from_slice(&self.read_chunk_scratch[..written]);
// Extend from the slice itself, not `.iter().copied()`: `VecDeque` specializes
// `Extend<&T> for T: Copy` into a bulk `copy_slice`, while the by-value iterator
// falls back to element-wise writes.
self.buffer.extend(&self.read_chunk_scratch[..written]);
}
let available = self.buffer.len() - self.buffer_cursor;
let available_samples = available / self.sample_size;
let available_samples = self.buffer.len() / self.sample_size;

if available_samples == 0 {
return Ok(None);
}

let batch_n = available_samples.min(self.batch_size);
let start = self.buffer_cursor;
let end = start + batch_n * self.sample_size;
self.buffer_cursor = end;
let take = batch_n * self.sample_size;
self.batches_yielded += 1;

let data = match recycled.and_then(T::from_recycled) {
let mut buf = match recycled.and_then(T::from_recycled) {
Some(mut buf) => {
buf.clear();
buf.extend_from_slice(&self.buffer[start..end]);
T::wrap(buf)
buf
}
None => T::wrap(self.buffer[start..end].to_vec()),
None => Vec::with_capacity(take),
};

if self.buffer_cursor >= self.buffer.len() / BUFFER_COMPACT_DENOM {
self.buffer.drain(..self.buffer_cursor);
self.buffer_cursor = 0;
// Copy the batch out of the ring's two slices, then drop the drain to discard the prefix.
// Copying via `as_slices` keeps the batch copy on `extend_from_slice`'s bulk path
// (`Drain` is not `TrustedLen`, so `extend(drain)` would copy element by element), and
// the front-anchored `drain(..take)` only advances the head — the retained tail is never
// shifted. Reusing a recycled Vec keeps the batch copy itself allocation-free.
let (front, back) = self.buffer.as_slices();
let n_front = take.min(front.len());
buf.extend_from_slice(&front[..n_front]);
if n_front < take {
buf.extend_from_slice(&back[..take - n_front]);
}
drop(self.buffer.drain(..take));
let data = T::wrap(buf);

Ok(Some(PrefetchedBatch {
data,
Expand Down Expand Up @@ -418,9 +426,6 @@ fn spawn_producer(
/// Default Parquet row group size for streaming reader (tunable).
const DEFAULT_PARQUET_ROW_GROUP_SIZE: usize = 2048;

/// When buffer_cursor >= buffer.len() / BUFFER_COMPACT_DENOM, compact by draining consumed prefix.
const BUFFER_COMPACT_DENOM: usize = 2;

/// Returns the path extension as lowercase ASCII (e.g. "parquet"), or None if missing/non-UTF8.
fn path_extension_lower(path: &Path) -> Option<String> {
path.extension()
Expand Down Expand Up @@ -631,10 +636,17 @@ where

buffer.truncate(written);
let read_chunk_scratch = vec![T::default(); initial_cap];
// Pre-reserve the streaming steady state so the ring never reallocates mid-run. Each batch
// drains up to `required` elements from the front while a refill appends up to one
// `initial_cap` chunk to the back, so live length peaks at `required + initial_cap`.
// `VecDeque::from` reuses the Vec's existing `initial_cap` capacity (O(1), no realloc), so
// this only tops it up by `required` — making front-advance realloc-free from batch 0.
let required = config.batch_size * sample_size;
let mut buffer = VecDeque::from(buffer);
buffer.reserve((required + initial_cap).saturating_sub(buffer.len()));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This eagerly reserves space for a full batch plus one refill chunk, even when the file only contains a partial batch. For large samples, this can substantially increase memory usage; for example, 20-qubit amplitude data with batch_size=64 and f64 grows from roughly 8 MiB to 520 MiB.

Could we avoid reserving the theoretical peak before EOF is known, and use checked arithmetic for the capacity calculation?

let producer = StreamingProducer::<T> {
reader,
buffer,
buffer_cursor: 0,
read_chunk_scratch,
sample_size,
batch_size: config.batch_size,
Expand Down Expand Up @@ -1553,13 +1565,14 @@ mod tests {
use std::fs;
use std::sync::Arc;

fn write_f32_parquet(path: &std::path::Path) {
// 8 samples, each 4 features — matches amplitude encoding with 2 qubits (2^2=4)
/// Writes `n_samples` f32 samples of 4 features each — matches amplitude encoding with
/// 2 qubits (2^2=4).
fn write_f32_parquet_n(path: &std::path::Path, n_samples: usize) {
let item_field = Arc::new(Field::new("item", DataType::Float32, true));
let list_field = Field::new("data", DataType::FixedSizeList(item_field, 4), true);
let schema = Arc::new(Schema::new(vec![list_field]));
let mut builder = FixedSizeListBuilder::new(Float32Builder::new(), 4);
for _ in 0..8 {
for _ in 0..n_samples {
builder.values().append_slice(&[0.25_f32, 0.5, 0.75, 1.0]);
builder.append(true);
}
Expand All @@ -1571,6 +1584,10 @@ mod tests {
writer.close().unwrap();
}

fn write_f32_parquet(path: &std::path::Path) {
write_f32_parquet_n(path, 8);
}

static FILE_COUNTER: std::sync::atomic::AtomicUsize =
std::sync::atomic::AtomicUsize::new(0);

Expand Down Expand Up @@ -1776,8 +1793,7 @@ mod tests {
let scratch = vec![0.0_f32; CAP];
let mut producer = StreamingProducer::<f32> {
reader,
buffer,
buffer_cursor: 0,
buffer: VecDeque::from(buffer),
read_chunk_scratch: scratch,
sample_size,
batch_size: 4,
Expand All @@ -1796,6 +1812,117 @@ mod tests {
);
}

/// #1436: the `VecDeque` ring buffer must not grow across a long stream. A deliberately
/// small `read_chunk_scratch` forces `produce()` to refill the ring on nearly every batch,
/// so the test exercises the real hot-path pattern — front-drain the consumed prefix while
/// new chunks are appended to the back. A correct ring reuses the freed front slots for the
/// appended tail, so capacity settles to a constant: front-advance is O(1) amortized with no
/// reallocation and no per-batch shift of the retained tail. The `refills_observed` assert
/// guards against a regression where the whole stream is slurped up front (which would make
/// the capacity check vacuous).
///
/// The same loop also pins down the wrap-around batch copy: once the ring wraps, a batch
/// can straddle the physical end of the allocation, so `produce()` must stitch it together
/// from both halves of `as_slices()`. Reaching that case needs a read chunk that is *not*
/// a multiple of the batch stride: if every batch drained the ring empty, `VecDeque` would
/// reset the head to 0 on each refill and the ring would never wrap (`split_observed`
/// asserts the straddling case is actually reached).
#[test]
fn test_streaming_producer_buffer_capacity_constant_over_100_batches() {
const BATCH_SIZE: usize = 4;
const SAMPLE_LEN: usize = 4; // 2 qubits, amplitude encoding
const BATCHES: usize = 100;
// Elements consumed per batch.
const STRIDE: usize = BATCH_SIZE * SAMPLE_LEN;
// One sample more than a batch: produce() refills on nearly every batch and always
// leaves a remainder behind, so the head walks around the ring instead of resetting.
const READ_CHUNK: usize = STRIDE + SAMPLE_LEN;
// Each sample written by write_f32_parquet_n, repeated across the batch.
const SAMPLE: [f32; SAMPLE_LEN] = [0.25, 0.5, 0.75, 1.0];

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Every sample in this fixture has the same values, the test cannot detect samples being reordered, duplicated, or skipped across the wrap boundary. Reversing the front and back slices still passes this test.

let path = temp_parquet_path("streaming_cap");
// A few extra samples so the 100th produce() is a full batch, not EOF.
write_f32_parquet_n(&path, BATCH_SIZE * (BATCHES + 5));

let run = (|| -> Result<()> {
let reader = ParquetStreamingReader::<f32>::new(
&path,
Some(DEFAULT_PARQUET_ROW_GROUP_SIZE),
NullHandling::FillZero,
)?;
// Exactly the steady-state peak the production build path reserves: one batch
// plus one refill chunk. Sized this tightly, the head wraps within the first few
// batches, so straddling batch copies are exercised early and often.
let mut producer = StreamingProducer::<f32> {
reader,
buffer: VecDeque::with_capacity(STRIDE + READ_CHUNK),
read_chunk_scratch: vec![0.0_f32; READ_CHUNK],
sample_size: SAMPLE_LEN,
batch_size: BATCH_SIZE,
num_qubits: 2,
batches_yielded: 0,
batch_limit: usize::MAX,
};
Comment on lines +1855 to +1864

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

IMO, This test constructs the deque with the expected steady-state capacity directly, so it does not verify the reserve logic in build_streaming_producer, right? The production reservation could be removed or broken while this test still passes.

Could the test create the producer through the prod initialization path and then verify that capacity remains stable for 100 batches?


let assert_batch_values = |batch: &PrefetchedBatch, i: usize| match &batch.data {
BatchData::F32(v) => {
assert_eq!(v.len(), STRIDE, "batch {i} has the wrong length");
for (j, x) in v.iter().enumerate() {
assert_eq!(
*x,
SAMPLE[j % SAMPLE_LEN],
"batch {i} element {j} corrupted",
);
}
}
other => panic!("batch {i} must be F32, got {other:?}"),
};

// Capacity after the first produce() is the steady-state baseline.
let first = producer.produce(None)?.expect("stream must yield a batch");
assert_batch_values(&first, 0);
let baseline_cap = producer.buffer.capacity();

let mut refills_observed = false;
let mut split_observed = false;
for i in 1..BATCHES {
let len_before = producer.buffer.len();
// The ring already holds a full batch that is split across the wrap boundary:
// this produce() needs no refill, so it must copy the batch out of both
// halves of as_slices().
let (front, back) = producer.buffer.as_slices();
if len_before >= STRIDE && !back.is_empty() && front.len() < STRIDE {
split_observed = true;
}

let batch = producer.produce(None)?;
let batch = batch.unwrap_or_else(|| panic!("batch {i} unexpectedly empty"));
assert_batch_values(&batch, i);
// A refill ran during this produce() iff the ring gained elements beyond the
// one batch (STRIDE) it just drained from the front.
if producer.buffer.len() + STRIDE > len_before {
refills_observed = true;
}
assert_eq!(
producer.buffer.capacity(),
baseline_cap,
"buffer capacity grew at batch {i}: front-advance must not reallocate",
);
}
assert!(
refills_observed,
"test must exercise back-appends (ring wrap-around), not a single up-front slurp",
);
assert!(
split_observed,
"test must exercise a batch straddling the ring's wrap boundary",
);
Ok(())
})();

let _ = fs::remove_file(&path);
run.unwrap();
}

/// Direct unit test of the shared f64→f32 narrowing helper used by every
/// non-Parquet format (Arrow IPC, NumPy, PyTorch, TensorFlow). Covers both the
/// variant switch and the documented ±Inf behavior for values outside f32 range.
Expand Down
Loading