diff --git a/src/flow/src/batching_mode/engine.rs b/src/flow/src/batching_mode/engine.rs index c3a21d73a5..4c26b6c4b3 100644 --- a/src/flow/src/batching_mode/engine.rs +++ b/src/flow/src/batching_mode/engine.rs @@ -769,13 +769,12 @@ impl BatchingEngine { let persistence = if let Some(factory) = &self.persistence_factory { let table_info = table.table_info(); let meta = &table_info.meta; - let effective_mode = if task.config.exact_sequence_range_required { - crate::IncrementalMode::SequenceRange - } else if task.config.batch_opts.experimental_enable_incremental_read - && task - .sequence_range_capable() - .await - .is_ok_and(|capable| capable) + let effective_mode = if task.config.exact_sequence_range_required + || (task.config.batch_opts.experimental_enable_incremental_read + && task + .sequence_range_capable() + .await + .is_ok_and(|capable| capable)) { crate::IncrementalMode::SequenceRange } else { diff --git a/src/flow/src/batching_mode/utils.rs b/src/flow/src/batching_mode/utils.rs index 66256e3618..c09fc1ff40 100644 --- a/src/flow/src/batching_mode/utils.rs +++ b/src/flow/src/batching_mode/utils.rs @@ -547,7 +547,7 @@ pub fn analyze_incremental_aggregate_plan( merge_columns.push(IncrementalAggregateMergeColumn { input_field_name: input_field_name.clone(), output_field_name, - merge_op: merge_op.clone(), + merge_op, }); } }