datafusion_distributed/worker/
worker_service.rs1use crate::protocol::LocalWorkerContext;
2use crate::worker::{SingleWriteMultiRead, WorkerSessionBuilder};
3use crate::{DefaultSessionBuilder, TaskData, TaskKey};
4use datafusion::common::DataFusionError;
5use datafusion::execution::runtime_env::RuntimeEnv;
6use moka::future::Cache;
7use std::borrow::Cow;
8use std::sync::Arc;
9use std::time::Duration;
10use url::Url;
11
12const TASK_CACHE_TTI: Duration = Duration::from_mins(10);
13
14pub(crate) type ResultTaskData = Result<TaskData, Arc<DataFusionError>>;
15pub(crate) type TaskDataEntries = Cache<TaskKey, Arc<SingleWriteMultiRead<ResultTaskData>>>;
16
17#[derive(Clone)]
18pub struct Worker {
19 pub(super) runtime: Arc<RuntimeEnv>,
20 pub(crate) task_data_entries: Arc<TaskDataEntries>,
24 pub(super) session_builder: Arc<dyn WorkerSessionBuilder + Send + Sync>,
25 pub(crate) max_message_size: Option<usize>,
26 pub(super) version: Cow<'static, str>,
27}
28
29impl Default for Worker {
30 fn default() -> Self {
31 let cache = Cache::builder().time_to_idle(TASK_CACHE_TTI).build();
32 Self {
33 runtime: Arc::new(RuntimeEnv::default()),
34 task_data_entries: Arc::new(cache),
35 session_builder: Arc::new(DefaultSessionBuilder),
36 max_message_size: Some(usize::MAX),
37 version: Cow::Borrowed(""),
38 }
39 }
40}
41
42impl Worker {
43 pub fn from_session_builder(
46 session_builder: impl WorkerSessionBuilder + Send + Sync + 'static,
47 ) -> Self {
48 Self {
49 session_builder: Arc::new(session_builder),
50 ..Default::default()
51 }
52 }
53
54 pub fn with_runtime_env(mut self, runtime_env: Arc<RuntimeEnv>) -> Self {
57 self.runtime = runtime_env;
58 self
59 }
60
61 pub fn with_max_message_size(mut self, size: usize) -> Self {
72 self.max_message_size = Some(size);
73 self
74 }
75
76 pub fn with_version(mut self, version: impl Into<Cow<'static, str>>) -> Self {
78 self.version = version.into();
79 self
80 }
81
82 pub fn version(&self) -> &str {
84 &self.version
85 }
86
87 pub fn to_local_worker_context(&self, self_url: Url) -> LocalWorkerContext {
92 LocalWorkerContext {
93 local_worker: self.clone(),
94 self_url,
95 }
96 }
97
98 #[cfg(any(test, feature = "integration"))]
100 pub async fn tasks_running(&self) -> usize {
101 self.task_data_entries.run_pending_tasks().await;
104 self.task_data_entries.entry_count() as usize
105 }
106}