Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
51 changes: 18 additions & 33 deletions crates/modelardb_server/src/cluster.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,15 +48,15 @@ pub(crate) struct Cluster {
/// The remote data folder that each node in the cluster should be synchronized with.
/// When a table is created, dropped, vacuumed, or truncated, it is done in the
/// remote data folder first.
remote_data_folder: DataFolder,
remote_data_folder: Arc<DataFolder>,
}

impl Cluster {
/// Create and register the cluster metadata tables and save the given `node` in the cluster.
/// It is assumed that `node` corresponds to the local system running `modelardbd`. If the
/// cluster metadata tables do not exist and could not be created or the node could not be
/// saved, return [`ModelarDbServerError`].
pub(crate) async fn try_new(node: Node, remote_data_folder: DataFolder) -> Result<Self> {
pub(crate) async fn try_new(node: Node, remote_data_folder: Arc<DataFolder>) -> Result<Self> {
remote_data_folder
.create_and_register_cluster_metadata_tables()
.await?;
Expand Down Expand Up @@ -499,24 +499,19 @@ mod test {

// Create a normal table in the remote data folder that should be retrieved and created
// in the local data folder and one that already exists.
create_normal_table(
"normal_table_1",
"column",
data_folders.local_data_folder.clone(),
)
.await;
create_normal_table("normal_table_1", "column", &data_folders.local_data_folder).await;

create_normal_table(
"normal_table_1",
"column",
data_folders.maybe_remote_data_folder.clone().unwrap(),
data_folders.maybe_remote_data_folder.as_ref().unwrap(),
)
.await;

create_normal_table(
"normal_table_2",
"column",
data_folders.maybe_remote_data_folder.clone().unwrap(),
data_folders.maybe_remote_data_folder.as_ref().unwrap(),
)
.await;

Expand All @@ -538,17 +533,12 @@ mod test {

// Create a normal table in the local data folder with the same name as a normal table in
// the remote data folder, but with a different schema.
create_normal_table(
NORMAL_TABLE_NAME,
"local",
data_folders.local_data_folder.clone(),
)
.await;
create_normal_table(NORMAL_TABLE_NAME, "local", &data_folders.local_data_folder).await;

create_normal_table(
NORMAL_TABLE_NAME,
"remote",
data_folders.maybe_remote_data_folder.clone().unwrap(),
data_folders.maybe_remote_data_folder.as_ref().unwrap(),
)
.await;

Expand Down Expand Up @@ -576,21 +566,21 @@ mod test {
create_time_series_table(
"time_series_table_1",
"field",
data_folders.local_data_folder.clone(),
&data_folders.local_data_folder,
)
.await;

create_time_series_table(
"time_series_table_1",
"field",
data_folders.maybe_remote_data_folder.clone().unwrap(),
data_folders.maybe_remote_data_folder.as_ref().unwrap(),
)
.await;

create_time_series_table(
"time_series_table_2",
"field",
data_folders.maybe_remote_data_folder.clone().unwrap(),
data_folders.maybe_remote_data_folder.as_ref().unwrap(),
)
.await;

Expand All @@ -615,14 +605,14 @@ mod test {
create_time_series_table(
TIME_SERIES_TABLE_NAME,
"local",
data_folders.local_data_folder.clone(),
&data_folders.local_data_folder,
)
.await;

create_time_series_table(
TIME_SERIES_TABLE_NAME,
"remote",
data_folders.maybe_remote_data_folder.clone().unwrap(),
data_folders.maybe_remote_data_folder.as_ref().unwrap(),
)
.await;

Expand All @@ -646,17 +636,12 @@ mod test {
let data_folders = context.data_folders.clone();

// Create tables in the local data folder that are not in the remote data folder.
create_normal_table(
NORMAL_TABLE_NAME,
"local",
data_folders.local_data_folder.clone(),
)
.await;
create_normal_table(NORMAL_TABLE_NAME, "local", &data_folders.local_data_folder).await;

create_time_series_table(
TIME_SERIES_TABLE_NAME,
"local",
data_folders.local_data_folder.clone(),
&data_folders.local_data_folder,
)
.await;

Expand All @@ -673,7 +658,7 @@ mod test {

/// Create a normal table named `table_name` with a single column named `column_name` in
/// `data_folder`.
async fn create_normal_table(table_name: &str, column_name: &str, data_folder: DataFolder) {
async fn create_normal_table(table_name: &str, column_name: &str, data_folder: &DataFolder) {
let schema = Schema::new(vec![Field::new(column_name, ArrowValue::DATA_TYPE, false)]);

data_folder
Expand All @@ -687,7 +672,7 @@ mod test {
async fn create_time_series_table(
table_name: &str,
column_name: &str,
data_folder: DataFolder,
data_folder: &DataFolder,
) {
let query_schema = Arc::new(Schema::new(vec![
Field::new("timestamp", ArrowTimestamp::DATA_TYPE, false),
Expand Down Expand Up @@ -724,7 +709,7 @@ mod test {
/// data folder in the context uses a local [`DataFolder`].
async fn create_context(local_temp_dir: &TempDir, remote_temp_dir: &TempDir) -> Arc<Context> {
let temp_dir_url = local_temp_dir.path().to_str().unwrap();
let local_data_folder = DataFolder::open_local_url(temp_dir_url).await.unwrap();
let local_data_folder = Arc::new(DataFolder::open_local_url(temp_dir_url).await.unwrap());

let edge_node = Node::new("edge".to_owned(), ServerMode::Edge);
let cluster = create_cluster_with_node(remote_temp_dir, edge_node).await;
Expand Down Expand Up @@ -813,7 +798,7 @@ mod test {
/// Create a [`Cluster`] that uses a local [`DataFolder`] for the remote data folder.
async fn create_cluster_with_node(temp_dir: &TempDir, node: Node) -> Cluster {
let temp_dir_url = temp_dir.path().to_str().unwrap();
let local_data_folder = DataFolder::open_local_url(temp_dir_url).await.unwrap();
let local_data_folder = Arc::new(DataFolder::open_local_url(temp_dir_url).await.unwrap());

Cluster::try_new(node, local_data_folder.clone())
.await
Expand Down
12 changes: 6 additions & 6 deletions crates/modelardb_server/src/configuration.rs
Original file line number Diff line number Diff line change
Expand Up @@ -195,7 +195,7 @@ pub struct ConfigurationManager {
/// The mode of the write-ahead log used to determine whether data is logged before ingestion.
wal_mode: WalMode,
/// The local data folder that stores the configuration file at the root.
local_data_folder: DataFolder,
local_data_folder: Arc<DataFolder>,
/// The configuration of the system. This is stored in a separate type to allow for easier
/// serialization and deserialization.
configuration: Configuration,
Expand All @@ -208,7 +208,7 @@ impl ConfigurationManager {
/// if the corresponding CLI flags or environment variables are set. If the configuration file
/// could not be read or created, [`ModelarDbServerError`] is returned.
pub async fn try_new(
local_data_folder: DataFolder,
local_data_folder: Arc<DataFolder>,
cluster_mode: ClusterMode,
args: &ServerArgs,
) -> Result<Self> {
Expand Down Expand Up @@ -589,7 +589,7 @@ mod tests {
async fn test_invalid_configuration_in_configuration_file() {
let temp_dir = tempfile::tempdir().unwrap();
let local_url = temp_dir.path().to_str().unwrap();
let local_data_folder = DataFolder::open_local_url(local_url).await.unwrap();
let local_data_folder = Arc::new(DataFolder::open_local_url(local_url).await.unwrap());

// Multiple threads per component are not supported.
let invalid_configuration = Configuration {
Expand Down Expand Up @@ -619,7 +619,7 @@ mod tests {
async fn test_invalid_toml_in_configuration_file() {
let temp_dir = tempfile::tempdir().unwrap();
let local_url = temp_dir.path().to_str().unwrap();
let local_data_folder = DataFolder::open_local_url(local_url).await.unwrap();
let local_data_folder = Arc::new(DataFolder::open_local_url(local_url).await.unwrap());

// Write invalid TOML to the configuration file.
let path = temp_dir.path().join(CONFIGURATION_FILE_NAME);
Expand Down Expand Up @@ -1000,11 +1000,11 @@ mod tests {
Arc<RwLock<ConfigurationManager>>,
) {
let local_url = temp_dir.path().to_str().unwrap();
let local_data_folder = DataFolder::open_local_url(local_url).await.unwrap();
let local_data_folder = Arc::new(DataFolder::open_local_url(local_url).await.unwrap());

let target_dir = tempfile::tempdir().unwrap();
let target_url = target_dir.path().to_str().unwrap();
let remote_data_folder = DataFolder::open_local_url(target_url).await.unwrap();
let remote_data_folder = Arc::new(DataFolder::open_local_url(target_url).await.unwrap());

let data_folders = DataFolders::new(
local_data_folder.clone(),
Expand Down
2 changes: 1 addition & 1 deletion crates/modelardb_server/src/context.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1161,7 +1161,7 @@ mod tests {
/// Create a simple [`Context`] that uses `temp_dir` as the local data folder and query data folder.
async fn create_context(temp_dir: &TempDir) -> Arc<Context> {
let temp_dir_url = temp_dir.path().to_str().unwrap();
let local_data_folder = DataFolder::open_local_url(temp_dir_url).await.unwrap();
let local_data_folder = Arc::new(DataFolder::open_local_url(temp_dir_url).await.unwrap());

Arc::new(
Context::try_new(
Expand Down
27 changes: 16 additions & 11 deletions crates/modelardb_server/src/data_folders.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,8 @@

//! Implementation of a struct that provides access to the local and remote data storage components.

use std::sync::Arc;

use modelardb_storage::data_folder::DataFolder;
use modelardb_types::types::{Node, ServerMode};

Expand All @@ -25,20 +27,20 @@ use crate::{ClusterMode, Result, ServerMode as ServerModeArg};
#[derive(Clone)]
pub struct DataFolders {
/// Folder for storing metadata and data in Apache Parquet files on the local file system.
pub local_data_folder: DataFolder,
pub local_data_folder: Arc<DataFolder>,
/// Folder for storing metadata and data in Apache Parquet files in a remote object store.
pub maybe_remote_data_folder: Option<DataFolder>,
pub maybe_remote_data_folder: Option<Arc<DataFolder>>,
/// Folder from which metadata and data in Apache Parquet files will be read during query
/// execution. It is equivalent to `local_data_folder` when deployed on the edge and
/// `remote_data_folder` when deployed in the cloud.
pub query_data_folder: DataFolder,
pub query_data_folder: Arc<DataFolder>,
}

impl DataFolders {
pub fn new(
local_data_folder: DataFolder,
maybe_remote_data_folder: Option<DataFolder>,
query_data_folder: DataFolder,
local_data_folder: Arc<DataFolder>,
maybe_remote_data_folder: Option<Arc<DataFolder>>,
query_data_folder: Arc<DataFolder>,
) -> Self {
Self {
local_data_folder,
Expand All @@ -65,7 +67,8 @@ impl DataFolders {
remote_data_folder: None,
..
} => {
let local_data_folder = DataFolder::open_local_url(local_data_folder).await?;
let local_data_folder =
Arc::new(DataFolder::open_local_url(local_data_folder).await?);

Ok((
ClusterMode::SingleNode,
Expand All @@ -79,9 +82,10 @@ impl DataFolders {
credentials,
} => {
let remote_data_folder =
DataFolder::open_remote_url(remote_data_folder, credentials).await?;
Arc::new(DataFolder::open_remote_url(remote_data_folder, credentials).await?);

let local_data_folder = DataFolder::open_local_url(local_data_folder).await?;
let local_data_folder =
Arc::new(DataFolder::open_local_url(local_data_folder).await?);

let node = Node::new(url_with_port, ServerMode::Edge);
let cluster = Cluster::try_new(node, remote_data_folder.clone()).await?;
Expand All @@ -102,9 +106,10 @@ impl DataFolders {
credentials,
} => {
let remote_data_folder =
DataFolder::open_remote_url(remote_data_folder, credentials).await?;
Arc::new(DataFolder::open_remote_url(remote_data_folder, credentials).await?);

let local_data_folder = DataFolder::open_local_url(local_data_folder).await?;
let local_data_folder =
Arc::new(DataFolder::open_local_url(local_data_folder).await?);

let node = Node::new(url_with_port, ServerMode::Cloud);
let cluster = Cluster::try_new(node, remote_data_folder.clone()).await?;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ pub(super) struct CompressedDataManager {
/// Component that transfers saved compressed data to the remote data folder when it is necessary.
pub(super) data_transfer: Arc<RwLock<Option<DataTransfer>>>,
/// Folder containing all compressed data managed by the [`StorageEngine`](crate::storage::StorageEngine).
pub(crate) local_data_folder: DataFolder,
pub(crate) local_data_folder: Arc<DataFolder>,
/// The compressed segments before they are saved to persistent storage. The key is the name of
/// the time series table the compressed segments represents data points for so the Apache Parquet
/// files can be partitioned by table.
Expand All @@ -63,7 +63,7 @@ impl CompressedDataManager {
pub(super) fn new(
data_storage_compactor: Arc<RwLock<DataStorageCompactor>>,
data_transfer: Arc<RwLock<Option<DataTransfer>>>,
local_data_folder: DataFolder,
local_data_folder: Arc<DataFolder>,
channels: Arc<Channels>,
memory_pool: Arc<MemoryPool>,
wal_mode: WalMode,
Expand Down Expand Up @@ -580,7 +580,7 @@ mod tests {

// Create a local data folder and save a single time series table to the Delta Lake.
let temp_dir_url = temp_dir.path().to_str().unwrap();
let local_data_folder = DataFolder::open_local_url(temp_dir_url).await.unwrap();
let local_data_folder = Arc::new(DataFolder::open_local_url(temp_dir_url).await.unwrap());

let time_series_table_metadata = table::time_series_table_metadata();
local_data_folder
Expand Down
14 changes: 9 additions & 5 deletions crates/modelardb_server/src/storage/data_storage_compactor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,8 @@
//! table by merging those small files into fewer larger ones and vacuuming the files left behind,
//! reducing storage use and query time.

use std::sync::Arc;

use dashmap::DashMap;
use modelardb_storage::data_folder::DataFolder;
use tracing::debug;
Expand All @@ -31,7 +33,7 @@ use crate::error::Result;
pub(super) struct DataStorageCompactor {
/// The data folder containing all compressed data managed by the
/// [`StorageEngine`](crate::storage::StorageEngine).
local_data_folder: DataFolder,
local_data_folder: Arc<DataFolder>,
/// The target size, in bytes, of the files produced when a table is optimized. A table is
/// compacted once its `estimated_compactable_size_in_bytes` reaches this size, so the same
/// value decides both when to compact and how large the optimized files are.
Expand All @@ -56,7 +58,7 @@ impl DataStorageCompactor {
/// so small files written before a restart are not forgotten. If the files in `local_data_folder`
/// could not be read, return [`ModelarDbServerError`](crate::error::ModelarDbServerError).
pub(super) async fn try_new(
local_data_folder: DataFolder,
local_data_folder: Arc<DataFolder>,
optimize_target_file_size_in_bytes: u64,
vacuum_retention_period_in_seconds: u64,
) -> Result<Self> {
Expand Down Expand Up @@ -292,10 +294,10 @@ mod tests {
}

/// Create a [`DataFolder`] in a local [`TempDir`] containing a single time series table.
async fn create_local_data_folder_with_table() -> (TempDir, DataFolder) {
async fn create_local_data_folder_with_table() -> (TempDir, Arc<DataFolder>) {
let temp_dir = tempfile::tempdir().unwrap();
let temp_dir_url = temp_dir.path().to_str().unwrap();
let local_data_folder = DataFolder::open_local_url(temp_dir_url).await.unwrap();
let local_data_folder = Arc::new(DataFolder::open_local_url(temp_dir_url).await.unwrap());

let time_series_table_metadata = table::time_series_table_metadata();
local_data_folder
Expand Down Expand Up @@ -333,7 +335,9 @@ mod tests {
}

/// Create a [`DataStorageCompactor`] that compacts the tables in `local_data_folder`.
async fn create_data_storage_compactor(local_data_folder: DataFolder) -> DataStorageCompactor {
async fn create_data_storage_compactor(
local_data_folder: Arc<DataFolder>,
) -> DataStorageCompactor {
DataStorageCompactor::try_new(
local_data_folder,
OPTIMIZE_TARGET_FILE_SIZE_IN_BYTES,
Expand Down
Loading
Loading