#[cfg(not(feature = "tombstones"))]
use std::sync::Arc;
use super::reader::{self, RowScanStream};
#[cfg(not(feature = "tombstones"))]
use crate::storage::producer_fault::{FaultScope, ScanTaskSite};
use crate::types::ScanRow;
#[cfg(not(feature = "tombstones"))]
use crate::types::TableId;
use crate::RowKey;
#[cfg(not(feature = "tombstones"))]
pub(super) fn spawn_fanout_merge(
readers: Vec<Arc<reader::SSTableReader>>,
table_id: TableId,
start_key: Option<RowKey>,
end_key: Option<RowKey>,
schema: Option<crate::schema::TableSchema>,
buffer_size: usize,
) -> RowScanStream {
let (out_tx, out_rx) = tokio::sync::mpsc::channel(buffer_size.max(1));
let fault_scope = FaultScope::capture(|| {
readers
.first()
.map(|reader| reader.file_path())
.unwrap_or_default()
});
let meter = crate::observability::read_metrics::ReadOpMeter::start(None);
let task = tokio::spawn(async move {
fault_scope.checkpoint(ScanTaskSite::FanoutMerge);
let _admission = reader::scan_stream_windowed::scan_admission::admit().await;
let mut streams: Vec<RowScanStream> = readers
.into_iter()
.map(|reader| {
reader.scan_stream_admitted(
table_id.clone(),
start_key.clone(),
end_key.clone(),
schema.clone(),
buffer_size,
reader::scan_stream_windowed::scan_admission::ScanAdmission::Exempt,
reader::ScanErrorReporting::Nested,
)
})
.collect();
let token_of =
|key: &RowKey| crate::util::cassandra_murmur3::cassandra_murmur3_token(key.as_bytes());
let mut heads: Vec<Option<(i64, RowKey, ScanRow)>> = Vec::with_capacity(streams.len());
for stream in streams.iter_mut() {
match stream.recv().await {
Some(Ok((key, row))) => heads.push(Some((token_of(&key), key, row))),
Some(Err(e)) => {
let _ = out_tx.send(Err(e)).await;
return;
}
None => heads.push(None),
}
}
loop {
let mut min_idx: Option<usize> = None;
for (i, head) in heads.iter().enumerate() {
if let Some((ref token, ref key, _)) = head {
match min_idx {
None => min_idx = Some(i),
Some(m) => {
if let Some((ref min_token, ref min_key, _)) = heads[m] {
if (token, key) < (min_token, min_key) {
min_idx = Some(i);
}
}
}
}
}
}
let idx = match min_idx {
Some(idx) => idx,
None => break, };
let entry = match heads[idx].take() {
Some((_, key, row)) => (key, row),
None => break, };
match streams[idx].recv().await {
Some(Ok((key, row))) => heads[idx] = Some((token_of(&key), key, row)),
Some(Err(e)) => {
let _ = out_tx.send(Err(e)).await;
return;
}
None => {} }
if out_tx.send(Ok(entry)).await.is_err() {
return; }
}
});
RowScanStream::new_measured_rows(out_rx, task, meter)
}
pub(super) fn rechunk_into_batches(
mut per_row: RowScanStream,
buffer_size: usize,
) -> reader::BatchedScanStream {
use reader::scan_stream_windowed::BATCH_EMIT_ROWS;
let cap = buffer_size.div_ceil(BATCH_EMIT_ROWS).max(1);
let (tx, rx) = tokio::sync::mpsc::channel(cap);
let task = tokio::spawn(async move {
let mut batch: Vec<(RowKey, ScanRow)> = Vec::with_capacity(BATCH_EMIT_ROWS);
while let Some(item) = per_row.recv().await {
match item {
Ok(entry) => {
batch.push(entry);
if batch.len() >= BATCH_EMIT_ROWS {
if tx.send(Ok(std::mem::take(&mut batch))).await.is_err() {
return; }
batch.reserve(BATCH_EMIT_ROWS);
}
}
Err(e) => {
if !batch.is_empty() {
let _ = tx.send(Ok(std::mem::take(&mut batch))).await;
}
let _ = tx.send(Err(e)).await;
return;
}
}
}
if !batch.is_empty() {
let _ = tx.send(Ok(batch)).await;
}
});
reader::BatchedScanStream::new_over_counted_source(rx, task)
}
#[cfg(all(test, feature = "write-support", not(feature = "tombstones")))]
#[path = "scan_stream_fanout_panic_tests.rs"]
mod panic_tests;