|
|
@ -80,8 +80,6 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer |
|
|
private long pollInterval; |
|
|
private long pollInterval; |
|
|
@Value("${queue.calculated_fields.pack_processing_timeout:60000}") |
|
|
@Value("${queue.calculated_fields.pack_processing_timeout:60000}") |
|
|
private long packProcessingTimeout; |
|
|
private long packProcessingTimeout; |
|
|
@Value("${queue.calculated_fields.pool_size:8}") |
|
|
|
|
|
private int poolSize; |
|
|
|
|
|
|
|
|
|
|
|
private final TbRuleEngineQueueFactory queueFactory; |
|
|
private final TbRuleEngineQueueFactory queueFactory; |
|
|
private final CalculatedFieldStateService stateService; |
|
|
private final CalculatedFieldStateService stateService; |
|
|
@ -148,14 +146,6 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private boolean[] partitionsToBooleanIndexArray(Set<TopicPartitionInfo> partitions) { |
|
|
|
|
|
boolean[] myPartitions = new boolean[partitionService.getTotalCalculatedFieldPartitions()]; |
|
|
|
|
|
for (var tpi : partitions) { |
|
|
|
|
|
tpi.getPartition().ifPresent(partition -> myPartitions[partition] = true); |
|
|
|
|
|
} |
|
|
|
|
|
return myPartitions; |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private void processMsgs(List<TbProtoQueueMsg<ToCalculatedFieldMsg>> msgs, TbQueueConsumer<TbProtoQueueMsg<ToCalculatedFieldMsg>> consumer, QueueConfig config) throws Exception { |
|
|
private void processMsgs(List<TbProtoQueueMsg<ToCalculatedFieldMsg>> msgs, TbQueueConsumer<TbProtoQueueMsg<ToCalculatedFieldMsg>> consumer, QueueConfig config) throws Exception { |
|
|
List<IdMsgPair<ToCalculatedFieldMsg>> orderedMsgList = msgs.stream().map(msg -> new IdMsgPair<>(UUID.randomUUID(), msg)).toList(); |
|
|
List<IdMsgPair<ToCalculatedFieldMsg>> orderedMsgList = msgs.stream().map(msg -> new IdMsgPair<>(UUID.randomUUID(), msg)).toList(); |
|
|
ConcurrentMap<UUID, TbProtoQueueMsg<ToCalculatedFieldMsg>> pendingMap = orderedMsgList.stream().collect( |
|
|
ConcurrentMap<UUID, TbProtoQueueMsg<ToCalculatedFieldMsg>> pendingMap = orderedMsgList.stream().collect( |
|
|
|