use crate::execution_engine::QueryStageExecutor;
use ballista_core::JobId;
use log::debug;
use std::{fmt::Display, sync::Arc};
pub trait ExecutorMetricsCollector: Send + Sync {
fn record_stage(
&self,
job_id: &JobId,
stage_id: usize,
partition: usize,
plan: Arc<dyn QueryStageExecutor>,
);
}
#[derive(Default)]
pub struct LoggingMetricsCollector {}
impl ExecutorMetricsCollector for LoggingMetricsCollector {
fn record_stage(
&self,
job_id: &JobId,
stage_id: usize,
partition: usize,
plan: Arc<dyn QueryStageExecutor>,
) {
debug!(
"\n=== [{job_id}/{stage_id}/{partition}] Physical plan with metrics ===\n{plan}\n"
);
}
}
#[derive(Clone, Copy, Debug, serde::Deserialize, Default)]
#[cfg_attr(feature = "build-binary", derive(clap::ValueEnum))]
pub enum ExecutorMetricCollectionPolicy {
#[cfg_attr(feature = "build-binary", clap(name = "sys"))]
SystemOnly,
#[cfg_attr(feature = "build-binary", clap(name = "proc"))]
#[default]
ProcessOnly,
#[cfg_attr(feature = "build-binary", clap(name = "all"))]
SystemAndProcess,
Off,
}
impl Display for ExecutorMetricCollectionPolicy {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
ExecutorMetricCollectionPolicy::SystemOnly => f.write_str("sys"),
ExecutorMetricCollectionPolicy::ProcessOnly => f.write_str("proc"),
ExecutorMetricCollectionPolicy::SystemAndProcess => f.write_str("all"),
ExecutorMetricCollectionPolicy::Off => f.write_str("off"),
}
}
}
#[cfg(feature = "build-binary")]
impl std::str::FromStr for ExecutorMetricCollectionPolicy {
type Err = String;
fn from_str(s: &str) -> std::result::Result<Self, Self::Err> {
clap::ValueEnum::from_str(s, true)
}
}