diff --git a/Cargo.lock b/Cargo.lock index e8a62ac28..69c950edd 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1017,6 +1017,7 @@ dependencies = [ "async-trait", "clap", "clp-rust-utils", + "const_format", "non-empty-string", "rmp-serde", "serde", diff --git a/components/clp-rust-utils/src/clp_config/package/config.rs b/components/clp-rust-utils/src/clp_config/package/config.rs index a2540a762..9f2a481f5 100644 --- a/components/clp-rust-utils/src/clp_config/package/config.rs +++ b/components/clp-rust-utils/src/clp_config/package/config.rs @@ -169,6 +169,18 @@ impl Database { ) } + /// # Returns + /// + /// The column-metadata table name (`_column_metadata`). + #[must_use] + pub fn column_metadata_table_name(&self, dataset: Option<&str>) -> String { + format!( + "{}{}_column_metadata", + self.table_prefix, + resolve_dataset_name(dataset) + ) + } + /// # Returns /// /// The datasets table name `datasets`. diff --git a/components/compression-coordinator/Cargo.toml b/components/compression-coordinator/Cargo.toml index 134d11a53..be6952f98 100644 --- a/components/compression-coordinator/Cargo.toml +++ b/components/compression-coordinator/Cargo.toml @@ -16,6 +16,8 @@ anyhow = "1.0.100" async-trait = "0.1.89" clap = { version = "4.6.4", features = ["derive"] } clp-rust-utils = { path = "../clp-rust-utils" } +const_format = "0.2.35" +non-empty-string = "0.2.6" rmp-serde = "1.3.1" serde = { version = "1.0.228", features = ["derive"] } spider-client = { git = "https://github.com/y-scope/spider.git", branch = "main" } @@ -25,6 +27,3 @@ strsim = "0.11.1" thiserror = "2.0.18" tokio = { version = "1.52.3", features = ["time", "rt-multi-thread"] } tracing = "0.1.44" - -[dev-dependencies] -non-empty-string = "0.2.6" diff --git a/components/compression-coordinator/src/error.rs b/components/compression-coordinator/src/error.rs index 129503bbf..f8814abcb 100644 --- a/components/compression-coordinator/src/error.rs +++ b/components/compression-coordinator/src/error.rs @@ -1,11 +1,46 @@ //! The crate-level error type for the compression coordinator. +use clp_rust_utils::{job_config::ingestion::JobId as IngestionJobId, s3::S3ObjectMetadataId}; + /// Errors returned by the compression coordinator. #[derive(Debug, thiserror::Error)] pub enum Error { + #[error( + "duplicate S3 object metadata IDs {ids:?} requested for ingestion job {ingestion_job_id}" + )] + DuplicateS3ObjectMetadata { + ingestion_job_id: IngestionJobId, + ids: Vec, + }, + + #[error("S3 object metadata {id} has an empty `{field}`")] + EmptyS3ObjectMetadataField { + id: S3ObjectMetadataId, + field: &'static str, + }, + #[error("invalid dataset: {0}")] InvalidDataset(String), + #[error("failed to create metadata table `{table}`: {source}")] + MetadataTableCreation { + table: String, + #[source] + source: sqlx::Error, + }, + + #[error("missing S3 object metadata {id} for ingestion job {ingestion_job_id}")] + MissingS3ObjectMetadata { + ingestion_job_id: IngestionJobId, + id: S3ObjectMetadataId, + }, + + #[error("no S3 object metadata was requested for ingestion job {0}")] + NoS3ObjectMetadata(IngestionJobId), + + #[error("no S3 objects were partitioned into compression task inputs")] + NoTaskInputs, + #[error("S3 bucket mismatch: expected `{0}`, but got `{1}`")] S3BucketMismatch(String, String), @@ -26,6 +61,9 @@ pub enum Error { #[error("failed to serialize a task input: {0}")] TaskInputSerialization(#[from] rmp_serde::encode::Error), + #[error("number of compression tasks {0} exceeds `i32::MAX`")] + TooManyCompressionTasks(usize), + #[error("unsupported input config")] UnsupportedInputConfig, } diff --git a/components/compression-coordinator/src/job_handle.rs b/components/compression-coordinator/src/job_handle.rs index 8d47da8a4..e3651eb44 100644 --- a/components/compression-coordinator/src/job_handle.rs +++ b/components/compression-coordinator/src/job_handle.rs @@ -6,15 +6,22 @@ use clp_rust_utils::{ clp_config::package::config::Database, dataset::VALID_DATASET_NAME_REGEX, job_config::{ClpIoConfig, CompressionJobId, CompressionJobStatus, InputConfig}, + s3::{ObjectMetadata, S3ObjectMetadataId}, task_io::compression::{ClpSCompressionOption, S3InputSource}, }; +use const_format::formatcp; +use non_empty_string::NonEmptyString; use spider_core::{ - task::ExecutionPolicy, + task::{ExecutionPolicy, TimeoutPolicy}, types::id::{JobId as SpiderJobId, ResourceGroupId}, }; use sqlx::MySqlPool; -use crate::{Error, compression_job_submitter::S3CompressionJobSubmitter}; +use crate::{ + Error, + compression_job_submitter::{CompressionJobOutcome, S3CompressionJobSubmitter}, + partition::CompressionInputBuilder, +}; /// Options for a compression job running in Spider. pub struct SpiderOption { @@ -30,16 +37,16 @@ pub struct SpiderOption { /// /// * `SubmitterType` - The type of the job submitter for Spider job submission. pub struct S3CompressionJobHandle { - _db_pool: MySqlPool, - _db_config: Database, + db_pool: MySqlPool, + db_config: Database, compression_job_id: CompressionJobId, job_submitter: SubmitterType, resource_group_id: ResourceGroupId, - _input_config: InputConfig, + input_config: InputConfig, clp_s_compression_option: ClpSCompressionOption, dataset: Option, - _target_archive_size: u64, + target_archive_size: u64, spider_option: Arc, } @@ -95,15 +102,15 @@ impl S3CompressionJobHandle S3CompressionJobHandle S3CompressionJobHandle Result<(), Error> { let input_sources = self.prepare_task_inputs().await?; + let num_tasks = input_sources.len(); self.upsert_metadata_tables().await?; @@ -195,7 +202,7 @@ impl S3CompressionJobHandle S3CompressionJobHandle S3CompressionJobHandle Result, Error> { - todo!("implement me!") + // NOTE: We batch the metadata fetch to avoid hitting the maximum placeholder limit of + // MySQL. + const FETCH_CHUNK_SIZE: usize = 1000; + + let InputConfig::S3ObjectMetadataInputConfig { config } = &self.input_config else { + unreachable!("the input config is validated in the factory") + }; + if config.s3_object_metadata_ids.is_empty() { + return Err(Error::NoS3ObjectMetadata(config.ingestion_job_id)); + } + + let mut sorted_metadata_ids = config.s3_object_metadata_ids.clone(); + sorted_metadata_ids.sort_unstable(); + let mut duplicate_ids: Vec = Vec::new(); + sorted_metadata_ids.dedup_by(|a, b| { + if *a != *b { + return false; + } + if duplicate_ids.last().copied() != Some(*b) { + duplicate_ids.push(*b); + } + true + }); + if !duplicate_ids.is_empty() { + return Err(Error::DuplicateS3ObjectMetadata { + ingestion_job_id: config.ingestion_job_id, + ids: duplicate_ids, + }); + } + + let mut input_builder = CompressionInputBuilder::from_s3_config( + config.s3_config.clone(), + self.target_archive_size, + ); + for chunk in sorted_metadata_ids.chunks(FETCH_CHUNK_SIZE) { + let mut query_builder = sqlx::QueryBuilder::::new(format!( + "SELECT `id`, `bucket`, `key`, `size` FROM \ + `{INGESTED_S3_OBJECT_METADATA_TABLE_NAME}` WHERE `id` IN (" + )); + let mut separated_ids = query_builder.separated(", "); + for metadata_id in chunk { + separated_ids.push_bind(metadata_id); + } + query_builder + .push(") AND `ingestion_job_id` = ") + .push_bind(config.ingestion_job_id); + query_builder.push(" ORDER BY `id` ASC"); + + let metadata_rows = query_builder + .build_query_as::() + .fetch_all(&self.db_pool) + .await?; + + let mut row_iter = metadata_rows.into_iter().peekable(); + for &expected_id in chunk { + if row_iter.peek().is_none_or(|row| row.id != expected_id) { + return Err(Error::MissingS3ObjectMetadata { + ingestion_job_id: config.ingestion_job_id, + id: expected_id, + }); + } + let row = row_iter + .next() + .expect("row should have been validated by the previous peek"); + let bucket = NonEmptyString::new(row.bucket).map_err(|_| { + Error::EmptyS3ObjectMetadataField { + id: row.id, + field: "bucket", + } + })?; + let key = NonEmptyString::new(row.key).map_err(|_| { + Error::EmptyS3ObjectMetadataField { + id: row.id, + field: "key", + } + })?; + input_builder.add(ObjectMetadata { + bucket, + key, + size: row.size, + })?; + } + } + + let input_sources = input_builder.into_task_input_sources(); + if input_sources.is_empty() { + return Err(Error::NoTaskInputs); + } + + Ok(input_sources + .into_iter() + .map(|input_source| { + // NOTE: The timeout policy scales with the number of objects assigned to the task. + // Each object contributes three minutes to the soft timeout and five minutes to + // the hard timeout. This heuristic does not account for object size and can be + // refined in the future. + let num_objects = input_source.object_keys.len() as u64; + let timeout_policy = TimeoutPolicy { + soft_timeout_ms: num_objects * 3 * 60 * 1000, + hard_timeout_ms: num_objects * 5 * 60 * 1000, + }; + let execution_policy = ExecutionPolicy { + max_num_retry: self.spider_option.compression_task_max_retry, + timeout_policy, + ..ExecutionPolicy::default() + }; + (input_source, execution_policy) + }) + .collect()) } /// Ensures that the required metadata tables exist for the configured dataset. @@ -265,9 +386,50 @@ impl S3CompressionJobHandle Result<(), Error> { - todo!("implement me!") + let archives_table = self.db_config.archives_table_name(self.dataset.as_deref()); + let column_metadata_table = self + .db_config + .column_metadata_table_name(self.dataset.as_deref()); + + sqlx::query(&format!( + "CREATE TABLE IF NOT EXISTS `{archives_table}` ( + `pagination_id` BIGINT unsigned NOT NULL AUTO_INCREMENT, + `id` VARCHAR(64) NOT NULL, + `begin_timestamp` BIGINT NOT NULL, + `end_timestamp` BIGINT NOT NULL, + `uncompressed_size` BIGINT NOT NULL, + `size` BIGINT NOT NULL, + `creator_id` VARCHAR(64) NOT NULL, + `creation_ix` INT NOT NULL, + KEY `archives_creation_order` (`creator_id`,`creation_ix`) USING BTREE, + UNIQUE KEY `archive_id` (`id`) USING BTREE, + PRIMARY KEY (`pagination_id`) + )" + )) + .execute(&self.db_pool) + .await + .map_err(|source| Error::MetadataTableCreation { + table: archives_table, + source, + })?; + + sqlx::query(&format!( + "CREATE TABLE IF NOT EXISTS `{column_metadata_table}` ( + `name` VARCHAR(512) NOT NULL, + `type` TINYINT NOT NULL, + PRIMARY KEY (`name`, `type`) + )" + )) + .execute(&self.db_pool) + .await + .map_err(|source| Error::MetadataTableCreation { + table: column_metadata_table, + source, + })?; + + Ok(()) } /// Persists the Spider job ID and marks the compression job as running. @@ -279,9 +441,27 @@ impl S3CompressionJobHandle Result<(), Error> { - todo!("implement me!") + /// * [`Error::TooManyCompressionTasks`] if `num_tasks` exceeds `i32`'s range. + /// * Forwards [`sqlx::query::Query::execute`]'s return values on failure. + async fn persist_spider_job_id( + &self, + spider_job_id: SpiderJobId, + num_tasks: usize, + ) -> Result<(), Error> { + let num_tasks = + i32::try_from(num_tasks).map_err(|_| Error::TooManyCompressionTasks(num_tasks))?; + sqlx::query(formatcp!( + "UPDATE `{COMPRESSION_JOB_TABLE_NAME}` SET `spider_id` = ?, `status` = ?, `num_tasks` \ + = ?, `start_time` = CURRENT_TIMESTAMP(3) WHERE `id` = ?" + )) + .bind(spider_job_id.get()) + .bind(CompressionJobStatus::Running) + .bind(num_tasks) + .bind(self.compression_job_id) + .execute(&self.db_pool) + .await?; + + Ok(()) } /// Waits for the associated Spider job to complete and finalizes the compression job. @@ -293,9 +473,81 @@ impl S3CompressionJobHandle Result<(), Error> { - todo!("implement me!") + /// * Forwards [`S3CompressionJobSubmitter::run_s3_compression_job_to_completion`]'s return + /// values on failure. + /// * Forwards [`Self::update_job_status`]'s return values on failure. + /// * Forwards [`Self::get_job_status`]'s return values on failure. + async fn to_completion(&self, spider_job_id: SpiderJobId) -> Result<(), Error> { + let outcome = self + .job_submitter + .run_s3_compression_job_to_completion( + spider_job_id, + self.spider_option.initial_poll_backoff, + self.spider_option.max_poll_backoff, + ) + .await?; + tracing::info!( + compression_job_id = % self.compression_job_id, + spider_job_id = % spider_job_id, + outcome = ? outcome, + "Compression job reached a terminal state.", + ); + + match outcome { + // The commit task records the successful CLP job status and publishes the archives in + // the same transaction. + CompressionJobOutcome::Succeeded => Ok(()), + CompressionJobOutcome::Failed { error_message } => { + if CompressionJobStatus::Succeeded == self.get_job_status().await? { + // NOTE: The commit task may successfully proceed but fail to report to Spider's + // control unit. In that case, the compression outcome has already committed to + // CLP DB, and thus the job should be considered `Succeeded`. + return Ok(()); + } + self.update_job_status( + CompressionJobStatus::Failed, + Some(format!( + "The Spider compression job failed: {error_message}" + )), + ) + .await + } + CompressionJobOutcome::Cancelled => { + if CompressionJobStatus::Succeeded == self.get_job_status().await? { + // NOTE: The commit task may successfully proceed but fail to report to Spider's + // control unit. In that case, the compression outcome has already committed to + // CLP DB, and thus the job should be considered `Succeeded`. + return Ok(()); + } + self.update_job_status( + CompressionJobStatus::Killed, + Some("The Spider compression job was cancelled.".to_owned()), + ) + .await + } + } + } + + /// Reads the current status of the compression job from the CLP database. + /// + /// # Returns + /// + /// The current [`CompressionJobStatus`] of the compression job on success. + /// + /// # Errors + /// + /// Returns an error if: + /// + /// * Forwards [`sqlx::query::QueryScalar::fetch_one`]'s return values on failure. + async fn get_job_status(&self) -> Result { + let job_status: CompressionJobStatus = sqlx::query_scalar(formatcp!( + "SELECT `status` FROM `{COMPRESSION_JOB_TABLE_NAME}` WHERE `id` = ?" + )) + .bind(self.compression_job_id) + .fetch_one(&self.db_pool) + .await?; + + Ok(job_status) } /// Updates the compression job status in the CLP database. @@ -304,12 +556,35 @@ impl S3CompressionJobHandle, + job_status: CompressionJobStatus, + status_message: Option, ) -> Result<(), Error> { - todo!("implement me!") + let status_message = status_message.as_ref().map_or("", String::as_str); + sqlx::query(formatcp!( + "UPDATE `{COMPRESSION_JOB_TABLE_NAME}` SET `status` = ?, `status_msg` = ? WHERE `id` \ + = ?" + )) + .bind(job_status) + .bind(status_message) + .bind(self.compression_job_id) + .execute(&self.db_pool) + .await?; + + Ok(()) } } + +const COMPRESSION_JOB_TABLE_NAME: &str = "compression_jobs"; +const INGESTED_S3_OBJECT_METADATA_TABLE_NAME: &str = "ingested_s3_object_metadata"; + +/// A projection of the columns read from an ingested S3 object metadata row. +#[derive(Debug, sqlx::FromRow)] +struct S3ObjectMetadataRow { + id: S3ObjectMetadataId, + bucket: String, + key: String, + size: u64, +} diff --git a/components/compression-coordinator/src/partition.rs b/components/compression-coordinator/src/partition.rs index 75f85707f..2281caf07 100644 --- a/components/compression-coordinator/src/partition.rs +++ b/components/compression-coordinator/src/partition.rs @@ -38,6 +38,16 @@ impl CompressionInputBuilder { InputConfig::S3ObjectMetadataInputConfig { config } => config.s3_config.clone(), }; + Self::from_s3_config(s3_config, target_archive_size) + } + + /// Creates an empty builder from the S3 input settings and target archive size. + /// + /// # Returns + /// + /// A newly created [`CompressionInputBuilder`] with an empty buffer. + #[must_use] + pub(crate) const fn from_s3_config(s3_config: S3Config, target_archive_size: u64) -> Self { Self { buffer: Vec::new(), partitioned_task_inputs: Vec::new(),