use std::borrow::BorrowMut;
use std::convert::TryFrom;
use std::ops::Deref;
use std::sync::Arc;
use std::time::Duration;
use crate::core::checkpoint::CheckpointHandle;
use crate::core::cluster::MetadataStorageType;
use crate::core::cluster::TaskResourceInfo;
use crate::core::env::{StreamApp, StreamExecutionEnvironment};
use crate::core::properties::{InnerSystemProperties, Properties, SystemProperties};
use crate::core::runtime::{ClusterDescriptor, ManagerStatus};
use crate::dag::metadata::DagMetadata;
use crate::dag::DagManager;
use crate::deployment::TResourceManager;
use crate::metrics::metric::Gauge;
use crate::metrics::register_gauge;
use crate::runtime::context::Context;
use crate::runtime::coordinator::checkpoint_manager::CheckpointManager;
use crate::runtime::coordinator::heart_beat_manager::HeartbeatResult;
use crate::runtime::coordinator::task_distribution::build_cluster_descriptor;
use crate::runtime::coordinator::web_server::web_launch;
use crate::storage::metadata::{
loop_read_cluster_descriptor, loop_save_cluster_descriptor, loop_update_application_status,
MetadataStorage,
};
use crate::utils::date_time::timestamp_str;
pub mod checkpoint_manager;
pub mod heart_beat_manager;
pub mod task_distribution;
pub mod web_server;
pub(crate) struct CoordinatorTask<S, R>
where
S: StreamApp + 'static,
R: TResourceManager + 'static,
{
context: Arc<Context>,
stream_app: S,
metadata_storage_mode: MetadataStorageType,
resource_manager: R,
stream_env: StreamExecutionEnvironment,
startup_number: Gauge,
}
impl<S, R> CoordinatorTask<S, R>
where
S: StreamApp + 'static,
R: TResourceManager + 'static,
{
pub fn new(
context: Arc<Context>,
stream_app: S,
resource_manager: R,
stream_env: StreamExecutionEnvironment,
) -> Self {
let metadata_storage_mode = context.cluster_config.metadata_storage.clone();
let startup_number = register_gauge("startup_number", vec![]);
CoordinatorTask {
context,
stream_app,
metadata_storage_mode,
resource_manager,
stream_env,
startup_number,
}
}
pub fn run(&mut self) -> anyhow::Result<()> {
info!("coordinator start with mode {}", self.context.manager_type);
let application_properties = self.prepare_properties();
self.stream_app
.build_stream(&application_properties, self.stream_env.borrow_mut());
let dag_manager = {
let raw_stream_graph = self.stream_env.stream_manager.stream_graph.borrow();
DagManager::try_from(raw_stream_graph.deref())?
};
info!("DagManager build success");
let dag_metadata = DagMetadata::from(&dag_manager);
debug!("DagMetadata: {}", dag_metadata.to_string());
let mut cluster_descriptor = self.build_metadata(&dag_manager, &application_properties);
debug!("ApplicationDescriptor : {}", cluster_descriptor.to_string());
let ck_manager = self.build_checkpoint_manager(
&dag_metadata,
&application_properties,
cluster_descriptor.borrow_mut(),
);
ck_manager.run_align_task();
info!("start CheckpointManager align task");
self.web_serve(cluster_descriptor.borrow_mut(), ck_manager, dag_metadata);
info!(
"serve coordinator web ui {}",
&cluster_descriptor.coordinator_manager.web_address
);
self.resource_manager
.prepare(&self.context, &cluster_descriptor);
info!("ResourceManager prepared");
self.gauge_startup(&cluster_descriptor);
loop {
self.gauge_startup_number(cluster_descriptor.borrow_mut());
self.save_metadata(&cluster_descriptor);
info!("save metadata to storage");
self.stream_app.pre_worker_startup(&cluster_descriptor);
info!("pre-worker startup event");
let worker_task_ids = self.allocate_worker();
info!("allocate workers success");
self.waiting_worker_status_fine();
info!("all worker status is fine");
let heartbeat_result =
heart_beat_manager::start_heartbeat_timer(self.metadata_storage_mode.clone());
info!("heartbeat timer has interrupted");
self.stop_all_worker_tasks(worker_task_ids);
info!("stop all workers");
if let HeartbeatResult::End = heartbeat_result {
return Ok(());
}
}
}
fn prepare_properties(&self) -> Properties {
let mut application_properties = Properties::new();
application_properties.set_cluster_mode(self.context.cluster_mode);
self.stream_app
.prepare_properties(application_properties.borrow_mut());
let mut keys: Vec<&str> = application_properties
.as_map()
.keys()
.map(|x| x.as_str())
.collect();
keys.sort();
for key in keys {
let value = application_properties.as_map().get(key).unwrap();
info!("properties key={}, value={}", key, value,);
}
application_properties
}
fn build_metadata(
&mut self,
dag_manager: &DagManager,
application_properties: &Properties,
) -> ClusterDescriptor {
let cluster_descriptor =
build_cluster_descriptor(dag_manager, application_properties, &self.context);
cluster_descriptor
}
fn save_metadata(&self, cluster_descriptor: &ClusterDescriptor) {
let mut metadata_storage = MetadataStorage::new(&self.metadata_storage_mode);
loop_save_cluster_descriptor(metadata_storage.borrow_mut(), cluster_descriptor.clone());
}
fn build_checkpoint_manager(
&self,
dag_manager: &DagMetadata,
application_properties: &Properties,
cluster_descriptor: &mut ClusterDescriptor,
) -> CheckpointManager {
let checkpoint_ttl = application_properties
.get_checkpoint_ttl()
.unwrap_or_else(|_e| Duration::from_secs(1 * 60 * 60));
let mut ck_manager = CheckpointManager::new(
dag_manager,
&self.context,
cluster_descriptor,
checkpoint_ttl,
);
let operator_checkpoints = ck_manager.load().expect("load checkpoints error");
if operator_checkpoints.len() == 0 {
return ck_manager;
}
for task_manager_descriptor in &mut cluster_descriptor.worker_managers {
for task_descriptor in &mut task_manager_descriptor.task_descriptors {
let task_number = task_descriptor.task_id.task_number;
for operator in &mut task_descriptor.operators {
let cks = operator_checkpoints.get(&operator.operator_id).unwrap();
if cks.len() == 0 {
debug!("operator {:?} checkpoint not found", operator.operator_id);
continue;
}
let ck = cks
.iter()
.find(|ck| ck.task_id.task_number == task_number)
.unwrap();
operator.checkpoint_id = ck.checkpoint_id;
operator.checkpoint_handle = Some(CheckpointHandle {
handle: ck.handle.handle.clone(),
});
info!("operator {:?} checkpoint loaded", operator);
}
}
}
ck_manager
}
fn web_serve(
&self,
cluster_descriptor: &mut ClusterDescriptor,
checkpoint_manager: CheckpointManager,
dag_metadata: DagMetadata,
) {
let context = self.context.clone();
let metadata_storage_mode = self.metadata_storage_mode.clone();
let address = web_launch(
context,
metadata_storage_mode,
checkpoint_manager,
dag_metadata,
);
cluster_descriptor.coordinator_manager.web_address = address;
}
fn allocate_worker(&self) -> Vec<TaskResourceInfo> {
self.resource_manager
.worker_allocate(&self.stream_app, &self.stream_env)
.expect("try allocate worker error")
}
fn waiting_worker_status_fine(&self ) {
let mut metadata_storage = MetadataStorage::new(&self.metadata_storage_mode);
loop {
info!("waiting all workers status fine...");
let job_descriptor = loop_read_cluster_descriptor(&metadata_storage);
let unregister_worker = job_descriptor
.worker_managers
.iter()
.find(|x| x.status.ne(&ManagerStatus::Registered));
if unregister_worker.is_none() {
loop_update_application_status(
metadata_storage.borrow_mut(),
ManagerStatus::Registered,
);
info!("all workers status fine and Job update state to `Registered`");
job_descriptor.worker_managers.iter().for_each(|tm| {
info!(
"Registered List: `{}` registered at {}",
tm.task_manager_id,
timestamp_str(tm.latest_heart_beat_ts),
);
});
break;
}
std::thread::sleep(Duration::from_secs(3));
}
}
fn stop_all_worker_tasks(&self, worker_task_ids: Vec<TaskResourceInfo>) {
loop {
let rt = self.resource_manager.stop_workers(worker_task_ids.clone());
match rt {
Ok(_) => {
break;
}
Err(e) => {
error!("try stop all workers error. {}", e);
std::thread::sleep(Duration::from_secs(2));
}
}
}
}
fn gauge_startup(&self, cluster_descriptor: &ClusterDescriptor) {
let coordinator_manager = &cluster_descriptor.coordinator_manager;
let uptime = register_gauge("uptime", vec![]);
uptime.store(coordinator_manager.uptime as i64);
let v_cores = register_gauge("v_cores", vec![]);
v_cores.store(coordinator_manager.v_cores as i64);
let memory_mb = register_gauge("memory_mb", vec![]);
memory_mb.store(coordinator_manager.memory_mb as i64);
let num_task_managers = register_gauge("num_task_managers", vec![]);
num_task_managers.store(coordinator_manager.num_task_managers as i64);
}
fn gauge_startup_number(&self, cluster_descriptor: &mut ClusterDescriptor) {
cluster_descriptor.coordinator_manager.startup_number += 1;
self.startup_number
.store(cluster_descriptor.coordinator_manager.startup_number as i64);
}
}