13 changed files with 289 additions and 209 deletions
@ -0,0 +1,109 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.service.queue.consumer; |
||||
|
|
||||
|
import lombok.Builder; |
||||
|
import lombok.Getter; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
||||
|
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; |
||||
|
import org.thingsboard.server.queue.TbQueueConsumer; |
||||
|
import org.thingsboard.server.queue.TbQueueMsg; |
||||
|
|
||||
|
import java.util.List; |
||||
|
import java.util.Set; |
||||
|
import java.util.concurrent.ExecutorService; |
||||
|
import java.util.function.Supplier; |
||||
|
|
||||
|
@Slf4j |
||||
|
public class QueueConsumerManager<M extends TbQueueMsg> { |
||||
|
|
||||
|
private final String name; |
||||
|
private final MsgPackProcessor<M> msgPackProcessor; |
||||
|
private final long pollInterval; |
||||
|
private final ExecutorService consumerExecutor; |
||||
|
private final String threadPrefix; |
||||
|
|
||||
|
@Getter |
||||
|
private final TbQueueConsumer<M> consumer; |
||||
|
private volatile boolean stopped; |
||||
|
|
||||
|
@Builder |
||||
|
public QueueConsumerManager(String name, MsgPackProcessor<M> msgPackProcessor, |
||||
|
long pollInterval, Supplier<TbQueueConsumer<M>> consumerCreator, |
||||
|
ExecutorService consumerExecutor, String threadPrefix) { |
||||
|
this.name = name; |
||||
|
this.pollInterval = pollInterval; |
||||
|
this.msgPackProcessor = msgPackProcessor; |
||||
|
this.consumerExecutor = consumerExecutor; |
||||
|
this.threadPrefix = threadPrefix; |
||||
|
this.consumer = consumerCreator.get(); |
||||
|
} |
||||
|
|
||||
|
public void subscribe() { |
||||
|
consumer.subscribe(); |
||||
|
} |
||||
|
|
||||
|
public void subscribe(Set<TopicPartitionInfo> partitions) { |
||||
|
consumer.subscribe(partitions); |
||||
|
} |
||||
|
|
||||
|
public void launch() { |
||||
|
log.info("[{}] Launching consumer", name); |
||||
|
consumerExecutor.submit(() -> { |
||||
|
if (threadPrefix != null) { |
||||
|
ThingsBoardThreadFactory.addThreadNamePrefix(threadPrefix); |
||||
|
} |
||||
|
try { |
||||
|
consumerLoop(consumer); |
||||
|
} catch (Throwable e) { |
||||
|
log.error("Failure in consumer loop", e); |
||||
|
} |
||||
|
log.info("[{}] Consumer stopped", name); |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
private void consumerLoop(TbQueueConsumer<M> consumer) { |
||||
|
while (!stopped && !consumer.isStopped()) { |
||||
|
try { |
||||
|
List<M> msgs = consumer.poll(pollInterval); |
||||
|
if (msgs.isEmpty()) { |
||||
|
continue; |
||||
|
} |
||||
|
msgPackProcessor.process(msgs, consumer); |
||||
|
} catch (Exception e) { |
||||
|
if (!consumer.isStopped()) { |
||||
|
log.warn("Failed to process messages from queue", e); |
||||
|
try { |
||||
|
Thread.sleep(pollInterval); |
||||
|
} catch (InterruptedException interruptedException) { |
||||
|
log.trace("Failed to wait until the server has capacity to handle new requests", interruptedException); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public void stop() { |
||||
|
log.debug("[{}] Stopping consumer", name); |
||||
|
stopped = true; |
||||
|
consumer.unsubscribe(); |
||||
|
} |
||||
|
|
||||
|
public interface MsgPackProcessor<M extends TbQueueMsg> { |
||||
|
void process(List<M> msgs, TbQueueConsumer<M> consumer) throws Exception; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,40 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.queue.housekeeper; |
||||
|
|
||||
|
import lombok.Getter; |
||||
|
import org.springframework.beans.factory.annotation.Value; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.thingsboard.server.common.data.housekeeper.HousekeeperTaskType; |
||||
|
|
||||
|
import java.util.Set; |
||||
|
|
||||
|
@Component |
||||
|
@Getter |
||||
|
public class HousekeeperConfig { |
||||
|
|
||||
|
@Value("${queue.core.housekeeper.disabled-task-types:}") |
||||
|
private Set<HousekeeperTaskType> disabledTaskTypes; |
||||
|
@Value("${queue.core.housekeeper.task-processing-timeout-ms:120000}") |
||||
|
private int taskProcessingTimeout; |
||||
|
@Value("${queue.core.housekeeper.poll-interval-ms:500}") |
||||
|
private int pollInterval; |
||||
|
@Value("${queue.core.housekeeper.task-reprocessing-delay-ms:5000}") |
||||
|
private int taskReprocessingDelay; |
||||
|
@Value("${queue.core.housekeeper.max-reprocessing-attempts:10}") |
||||
|
private int maxReprocessingAttempts; |
||||
|
|
||||
|
} |
||||
Loading…
Reference in new issue