use std::ops::Range;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use bytes::Bytes;
use crate::error::Result;
use crate::io::FileRead;
use crate::scan::ArrowRecordBatchStream;
pub(crate) struct CountingFileRead<F: FileRead> {
inner: F,
bytes_read: Arc<AtomicU64>,
}
impl<F: FileRead> CountingFileRead<F> {
pub(crate) fn new(inner: F, bytes_read: Arc<AtomicU64>) -> Self {
Self { inner, bytes_read }
}
}
#[async_trait::async_trait]
impl<F: FileRead> FileRead for CountingFileRead<F> {
async fn read(&self, range: Range<u64>) -> Result<Bytes> {
debug_assert!(range.end >= range.start);
self.bytes_read
.fetch_add(range.end - range.start, Ordering::Relaxed);
self.inner.read(range).await
}
}
#[derive(Clone, Debug)]
pub struct ScanMetrics {
bytes_read: Arc<AtomicU64>,
}
impl ScanMetrics {
pub(crate) fn new() -> Self {
Self {
bytes_read: Arc::new(AtomicU64::new(0)),
}
}
pub(crate) fn bytes_read_counter(&self) -> &Arc<AtomicU64> {
&self.bytes_read
}
pub fn bytes_read(&self) -> u64 {
self.bytes_read.load(Ordering::Relaxed)
}
}
pub struct ScanResult {
stream: ArrowRecordBatchStream,
metrics: ScanMetrics,
}
impl ScanResult {
pub(crate) fn new(stream: ArrowRecordBatchStream, metrics: ScanMetrics) -> Self {
Self { stream, metrics }
}
pub fn stream(self) -> ArrowRecordBatchStream {
self.stream
}
pub fn metrics(&self) -> &ScanMetrics {
&self.metrics
}
}