Skip to content
Merged
Show file tree
Hide file tree
Changes from 2 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
6 changes: 5 additions & 1 deletion components/compression-coordinator/src/coordination.rs
Original file line number Diff line number Diff line change
Expand Up @@ -290,6 +290,9 @@ impl Coordinator {

/// Marks the compression jobs identified by `job_ids` with the current dispatch time.
///
/// If the `dispatch_time` has already been set by the `job_handler`, we preserve the value and
/// skip the update. See [`S3CompressionJobHandle::persist_spider_job_id`] for details.
///
Comment thread
Bill-hbrhbr marked this conversation as resolved.
Outdated
/// # Errors
///
/// Returns an error if:
Expand All @@ -305,7 +308,8 @@ impl Coordinator {
let mut tx = self.db_pool.begin().await?;
for chunk in job_ids.chunks(1000) {
let mut query_builder = sqlx::QueryBuilder::<sqlx::MySql>::new(formatcp!(
"UPDATE `{table}` SET `dispatch_time` = CURRENT_TIMESTAMP() WHERE `id` IN (",
"UPDATE `{table}` SET `dispatch_time` = COALESCE(`dispatch_time`, \
CURRENT_TIMESTAMP()) WHERE `id` IN (",
table = COMPRESSION_JOB_TABLE_NAME,
));
let mut separated_ids = query_builder.separated(", ");
Expand Down
9 changes: 8 additions & 1 deletion components/compression-coordinator/src/job_handle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -440,6 +440,12 @@ impl<SubmitterType: S3CompressionJobSubmitter> S3CompressionJobHandle<SubmitterT
/// This method associates the given Spider job ID with the compression job in the CLP database
/// and updates the compression job status to [`CompressionJobStatus::Running`].
///
/// This method also ensures that the job has a valid `dispatch_time`, which the coordinator
/// uses to mark jobs as dispatched. A coordinator restart may occur before the marker is
/// persisted, leaving the Spider job running without a valid `dispatch_time`. Therefore, this
/// method sets the field as part of row update if it has not already been set by the
/// coordinator.
///
/// # Errors
///
/// Returns an error if:
Expand All @@ -455,7 +461,8 @@ impl<SubmitterType: S3CompressionJobSubmitter> S3CompressionJobHandle<SubmitterT
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` = ?"
= ?, `start_time` = CURRENT_TIMESTAMP(3), `dispatch_time` = \
COALESCE(`dispatch_time`, CURRENT_TIMESTAMP()) WHERE `id` = ?"
))
.bind(spider_job_id.get())
.bind(CompressionJobStatus::Running)
Expand Down
Loading