#![doc = include_str!("../README.md")]
#![warn(missing_docs)]
mod client_pool;
pub mod collect;
#[cfg(feature = "build-binary")]
pub mod config;
pub mod execution_engine;
pub mod execution_loop;
pub mod executor;
pub mod executor_process;
pub mod executor_server;
pub mod flight_service;
pub mod metrics;
pub mod runtime_cache;
pub mod shutdown;
pub mod terminate;
mod cpu_bound_executor;
mod standalone;
use ballista_core::error::BallistaError;
use log::debug;
use std::net::SocketAddr;
pub use standalone::new_standalone_executor;
pub use standalone::new_standalone_executor_from_builder;
pub use standalone::new_standalone_executor_from_state;
use log::info;
use crate::shutdown::Shutdown;
use ballista_core::serde::protobuf::{
FailedTask, OperatorMetricsSet, ShuffleWritePartition, SuccessfulTask, TaskStatus,
task_status,
};
use ballista_core::serde::scheduler::PartitionId;
use ballista_core::utils::GrpcServerConfig;
pub type ArrowFlightServerProvider = dyn Fn(
String,
SocketAddr,
Shutdown,
GrpcServerConfig,
) -> tokio::task::JoinHandle<Result<(), BallistaError>>
+ Send
+ Sync;
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct TaskExecutionTimes {
launch_time: u64,
start_exec_time: u64,
end_exec_time: u64,
}
pub fn as_task_status(
execution_result: ballista_core::error::Result<Vec<ShuffleWritePartition>>,
executor_id: String,
task_id: usize,
stage_attempt_num: usize,
partition_id: PartitionId,
operator_metrics: Option<Vec<OperatorMetricsSet>>,
execution_times: TaskExecutionTimes,
) -> TaskStatus {
let metrics = operator_metrics.unwrap_or_default();
match execution_result {
Ok(partitions) => {
debug!(
"Task {:?} finished with operator_metrics array size {}",
task_id,
metrics.len()
);
TaskStatus {
task_id: task_id as u32,
job_id: partition_id.job_id.into(),
stage_id: partition_id.stage_id as u32,
stage_attempt_num: stage_attempt_num as u32,
partition_id: partition_id.partition_id as u32,
launch_time: execution_times.launch_time,
start_exec_time: execution_times.start_exec_time,
end_exec_time: execution_times.end_exec_time,
metrics,
status: Some(task_status::Status::Successful(SuccessfulTask {
executor_id,
partitions,
})),
}
}
Err(e) => {
let error_msg = e.to_string();
info!("Task {task_id:?} failed: {error_msg}");
TaskStatus {
task_id: task_id as u32,
job_id: partition_id.job_id.into(),
stage_id: partition_id.stage_id as u32,
stage_attempt_num: stage_attempt_num as u32,
partition_id: partition_id.partition_id as u32,
launch_time: execution_times.launch_time,
start_exec_time: execution_times.start_exec_time,
end_exec_time: execution_times.end_exec_time,
metrics,
status: Some(task_status::Status::Failed(FailedTask::from(e))),
}
}
}
}