Skip to content

Commit 2831196

Browse files
davsclausclaude
andauthored
CAMEL-24079: Fix QuartzScheduledPollConsumerScheduler ignoring startScheduler=false
Move quartzScheduler.scheduleJob()/rescheduleJob() from doStart() to startScheduler() so the startScheduler=false flag is respected. Same structural fix as CAMEL-24066 for the Spring scheduler. doStart() now only prepares resources; startScheduler() begins actual scheduling. Closes #24715 Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
1 parent c19b987 commit 2831196

2 files changed

Lines changed: 94 additions & 36 deletions

File tree

components/camel-quartz/src/main/java/org/apache/camel/pollconsumer/quartz/QuartzScheduledPollConsumerScheduler.java

Lines changed: 38 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,8 @@ public class QuartzScheduledPollConsumerScheduler extends ServiceSupport
7171
private boolean deleteJob = true;
7272
private volatile CronTrigger trigger;
7373
private volatile JobDetail job;
74+
private volatile boolean scheduled;
75+
private volatile TriggerKey rescheduleKey;
7476

7577
@Override
7678
public void onInit(Consumer consumer) {
@@ -100,16 +102,43 @@ public void unscheduleTask() {
100102

101103
@Override
102104
public void startScheduler() {
103-
// the quartz component starts the scheduler
105+
if (!scheduled) {
106+
try {
107+
if (rescheduleKey != null) {
108+
LOG.debug("Re-scheduling job: {} with trigger: {}", job, trigger.getKey());
109+
quartzScheduler.rescheduleJob(rescheduleKey, trigger);
110+
} else {
111+
LOG.debug("Scheduling job: {} with trigger: {}", job, trigger.getKey());
112+
try {
113+
quartzScheduler.scheduleJob(job, trigger);
114+
} catch (ObjectAlreadyExistsException ex) {
115+
QuartzComponent quartz = getCamelContext().getComponent("quartz", QuartzComponent.class);
116+
if (!(quartz.isClustered())) {
117+
throw ex;
118+
} else {
119+
TriggerKey triggerKey = trigger.getKey();
120+
trigger = (CronTrigger) quartzScheduler.getTrigger(triggerKey);
121+
if (trigger == null) {
122+
throw new SchedulerException("Trigger could not be found in quartz scheduler.");
123+
}
124+
}
125+
}
126+
}
127+
scheduled = true;
128+
if (LOG.isInfoEnabled()) {
129+
LOG.info("Job {} (triggerType={}, jobClass={}) is scheduled. Next fire date is {}",
130+
trigger.getKey(), trigger.getClass().getSimpleName(),
131+
job.getJobClass().getSimpleName(), trigger.getNextFireTime());
132+
}
133+
} catch (SchedulerException e) {
134+
throw RuntimeCamelException.wrapRuntimeCamelException(e);
135+
}
136+
}
104137
}
105138

106139
@Override
107140
public boolean isSchedulerStarted() {
108-
try {
109-
return quartzScheduler != null && quartzScheduler.isStarted();
110-
} catch (SchedulerException e) {
111-
return false;
112-
}
141+
return scheduled;
113142
}
114143

115144
@Override
@@ -244,9 +273,6 @@ protected void doStart() throws Exception {
244273
LOG.debug("Setting user extra triggerParameters {}", copy);
245274
PropertyBindingSupport.bindProperties(camelContext, trigger, copy);
246275
}
247-
248-
LOG.debug("Scheduling job: {} with trigger: {}", job, trigger.getKey());
249-
quartzScheduler.scheduleJob(job, trigger);
250276
} else {
251277
checkTriggerIsNonConflicting(existingTrigger);
252278

@@ -264,36 +290,10 @@ protected void doStart() throws Exception {
264290
.withSchedule(CronScheduleBuilder.cronSchedule(getCron()).inTimeZone(getTimeZone()))
265291
.build();
266292

267-
// Reschedule job if trigger settings were changed
268293
if (hasTriggerChanged(existingTrigger, trigger)) {
269-
LOG.debug("Re-scheduling job: {} with trigger: {}", job, trigger.getKey());
270-
quartzScheduler.rescheduleJob(triggerKey, trigger);
271-
} else {
272-
// Schedule it now. Remember that scheduler might not be started it, but we can schedule now.
273-
LOG.debug("Scheduling job: {} with trigger: {}", job, trigger.getKey());
274-
try {
275-
// Schedule it now. Remember that scheduler might not be started it, but we can schedule now.
276-
quartzScheduler.scheduleJob(job, trigger);
277-
} catch (ObjectAlreadyExistsException ex) {
278-
// some other VM might may have stored the job & trigger in DB in clustered mode, in the mean time
279-
QuartzComponent quartz = getCamelContext().getComponent("quartz", QuartzComponent.class);
280-
if (!(quartz.isClustered())) {
281-
throw ex;
282-
} else {
283-
trigger = (CronTrigger) quartzScheduler.getTrigger(triggerKey);
284-
if (trigger == null) {
285-
throw new SchedulerException("Trigger could not be found in quartz scheduler.");
286-
}
287-
}
288-
}
294+
rescheduleKey = triggerKey;
289295
}
290296
}
291-
292-
if (LOG.isInfoEnabled()) {
293-
LOG.info("Job {} (triggerType={}, jobClass={}) is scheduled. Next fire date is {}",
294-
trigger.getKey(), trigger.getClass().getSimpleName(),
295-
job.getJobClass().getSimpleName(), trigger.getNextFireTime());
296-
}
297297
}
298298

299299
@Override
@@ -313,6 +313,8 @@ private void unscheduleJob() throws SchedulerException {
313313
quartzScheduler.unscheduleJob(trigger.getKey());
314314
}
315315
}
316+
scheduled = false;
317+
rescheduleKey = null;
316318
}
317319

318320
private void checkTriggerIsNonConflicting(Trigger trigger) {
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,56 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one or more
3+
* contributor license agreements. See the NOTICE file distributed with
4+
* this work for additional information regarding copyright ownership.
5+
* The ASF licenses this file to You under the Apache License, Version 2.0
6+
* (the "License"); you may not use this file except in compliance with
7+
* the License. You may obtain a copy of the License at
8+
*
9+
* http://www.apache.org/licenses/LICENSE-2.0
10+
*
11+
* Unless required by applicable law or agreed to in writing, software
12+
* distributed under the License is distributed on an "AS IS" BASIS,
13+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14+
* See the License for the specific language governing permissions and
15+
* limitations under the License.
16+
*/
17+
package org.apache.camel.pollconsumer.quartz;
18+
19+
import java.nio.file.Path;
20+
import java.util.concurrent.TimeUnit;
21+
22+
import org.apache.camel.Exchange;
23+
import org.apache.camel.builder.RouteBuilder;
24+
import org.apache.camel.test.junit6.CamelTestSupport;
25+
import org.apache.camel.test.junit6.TestSupport;
26+
import org.junit.jupiter.api.Test;
27+
import org.junit.jupiter.api.io.TempDir;
28+
29+
import static org.awaitility.Awaitility.await;
30+
31+
public class FileConsumerQuartzSchedulerStartSchedulerTest extends CamelTestSupport {
32+
@TempDir
33+
Path testDirectory;
34+
35+
@Test
36+
public void testStartSchedulerFalseMustNotPoll() throws Exception {
37+
template.sendBodyAndHeader(TestSupport.fileUri(testDirectory), "Hello World", Exchange.FILE_NAME, "hello.txt");
38+
39+
getMockEndpoint("mock:result").expectedMessageCount(0);
40+
41+
// startScheduler=false should prevent polling even with scheduler=quartz
42+
await().during(3, TimeUnit.SECONDS).atMost(4, TimeUnit.SECONDS)
43+
.untilAsserted(() -> getMockEndpoint("mock:result").assertIsSatisfied());
44+
}
45+
46+
@Override
47+
protected RouteBuilder createRouteBuilder() {
48+
return new RouteBuilder() {
49+
@Override
50+
public void configure() {
51+
from(TestSupport.fileUri(testDirectory, "?scheduler=quartz&scheduler.cron=0/2+*+*+*+*+?&startScheduler=false"))
52+
.to("mock:result");
53+
}
54+
};
55+
}
56+
}

0 commit comments

Comments
 (0)