use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::task::{Context, Poll};
use futures::Stream;
use crate::exec::{FlowResult, ValueBatch, ValueBatchStream};
fn now_ns() -> u64 {
use web_time::Instant;
thread_local! {
static BASE: Instant = Instant::now();
}
BASE.with(|base| base.elapsed().as_nanos() as u64)
}
#[derive(Debug)]
pub(crate) struct OperatorMetrics {
enabled: AtomicBool,
output_rows: AtomicU64,
output_batches: AtomicU64,
elapsed_ns: AtomicU64,
skipped_rows: AtomicU64,
edges_scanned: AtomicU64,
}
impl OperatorMetrics {
pub(crate) fn new() -> Self {
Self {
enabled: AtomicBool::new(false),
output_rows: AtomicU64::new(0),
output_batches: AtomicU64::new(0),
elapsed_ns: AtomicU64::new(0),
skipped_rows: AtomicU64::new(0),
edges_scanned: AtomicU64::new(0),
}
}
pub(crate) fn is_enabled(&self) -> bool {
self.enabled.load(Ordering::Relaxed)
}
pub(crate) fn enable(&self) {
self.enabled.store(true, Ordering::Relaxed);
}
pub(crate) fn output_rows(&self) -> u64 {
self.output_rows.load(Ordering::Relaxed)
}
pub(crate) fn output_batches(&self) -> u64 {
self.output_batches.load(Ordering::Relaxed)
}
pub(crate) fn elapsed_ns(&self) -> u64 {
self.elapsed_ns.load(Ordering::Relaxed)
}
pub(crate) fn skipped_rows(&self) -> u64 {
self.skipped_rows.load(Ordering::Relaxed)
}
pub(crate) fn add_skipped_rows(&self, n: u64) {
self.skipped_rows.fetch_add(n, Ordering::Relaxed);
}
pub(crate) fn edges_scanned(&self) -> u64 {
self.edges_scanned.load(Ordering::Relaxed)
}
pub(crate) fn add_edges_scanned(&self, n: u64) {
self.edges_scanned.fetch_add(n, Ordering::Relaxed);
}
fn record_batch(&self, rows: u64, delta_ns: u64) {
self.output_rows.fetch_add(rows, Ordering::Relaxed);
self.output_batches.fetch_add(1, Ordering::Relaxed);
self.elapsed_ns.fetch_add(delta_ns, Ordering::Relaxed);
}
fn record_elapsed(&self, delta_ns: u64) {
self.elapsed_ns.fetch_add(delta_ns, Ordering::Relaxed);
}
}
struct MetricsStream {
inner: ValueBatchStream,
metrics: Arc<OperatorMetrics>,
name: &'static str,
batch_idx: u64,
}
impl Stream for MetricsStream {
type Item = FlowResult<ValueBatch>;
fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
let this = self.get_mut(); let span = tracing::trace_span!(
"batch",
op = this.name,
idx = this.batch_idx,
size = tracing::field::Empty,
);
let _enter = span.enter();
let start = now_ns();
let result = this.inner.as_mut().poll_next(cx);
let delta = now_ns().saturating_sub(start);
match &result {
Poll::Ready(Some(Ok(batch))) => {
let rows = batch.values.len() as u64;
span.record("size", rows);
this.metrics.record_batch(rows, delta);
this.batch_idx += 1;
}
Poll::Ready(Some(Err(_))) | Poll::Ready(None) => {
this.metrics.record_elapsed(delta);
}
Poll::Pending => {
this.metrics.record_elapsed(delta);
}
}
result
}
}
pub(crate) fn monitor_stream(
stream: ValueBatchStream,
name: &'static str,
metrics: &Arc<OperatorMetrics>,
) -> ValueBatchStream {
if !metrics.enabled.load(Ordering::Relaxed) {
return stream;
}
Box::pin(MetricsStream {
inner: stream,
metrics: Arc::clone(metrics),
name,
batch_idx: 0,
})
}