|
|
|
@ -433,16 +433,13 @@ public class TbRuleEngineQueueConsumerManager { |
|
|
|
removedPartitions.forEach((tpi) -> { |
|
|
|
consumers.get(tpi).initiateStop(); |
|
|
|
}); |
|
|
|
removedPartitions.forEach((tpi) -> { |
|
|
|
consumers.remove(tpi).awaitCompletion(); |
|
|
|
}); |
|
|
|
|
|
|
|
|
|
|
|
// ctx.getQueueAdmin().createTopicIfNotExists();
|
|
|
|
addedPartitions.forEach((tpi) -> { |
|
|
|
int partitionId = tpi.getPartition().orElse(-999999); |
|
|
|
String key = queueKey + "-" + partitionId; |
|
|
|
TbQueueConsumerTask consumer = new TbQueueConsumerTask(key, ctx.getQueueFactory().createToRuleEngineMsgConsumer(queue, null)); |
|
|
|
TbQueueConsumerTask consumer = new TbQueueConsumerTask(key, () -> ctx.getQueueFactory().createToRuleEngineMsgConsumer(queue, partitionId)); |
|
|
|
consumers.put(tpi, consumer); |
|
|
|
consumer.subscribe(Set.of(tpi)); |
|
|
|
launchConsumer(consumer); |
|
|
|
@ -471,7 +468,7 @@ public class TbRuleEngineQueueConsumerManager { |
|
|
|
} |
|
|
|
|
|
|
|
if (consumer == null) { |
|
|
|
consumer = new TbQueueConsumerTask(queueKey, ctx.getQueueFactory().createToRuleEngineMsgConsumer(queue, null)); |
|
|
|
consumer = new TbQueueConsumerTask(queueKey, () -> ctx.getQueueFactory().createToRuleEngineMsgConsumer(queue, null)); |
|
|
|
} |
|
|
|
consumer.subscribe(partitions); |
|
|
|
if (!consumer.isRunning()) { |
|
|
|
|