Browse Source

Merge pull request #9422 from thingsboard/fix/notification-rule-processor-blocking

DefaultNotificationRuleProcessor - submit to executor sooner
pull/9431/head
Andrew Shvayka 3 years ago
committed by GitHub
parent
commit
4155f67b04
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 31
      application/src/main/java/org/thingsboard/server/service/notification/rule/DefaultNotificationRuleProcessor.java

31
application/src/main/java/org/thingsboard/server/service/notification/rule/DefaultNotificationRuleProcessor.java

@ -77,18 +77,17 @@ public class DefaultNotificationRuleProcessor implements NotificationRuleProcess
public void process(NotificationRuleTrigger trigger) { public void process(NotificationRuleTrigger trigger) {
NotificationRuleTriggerType triggerType = trigger.getType(); NotificationRuleTriggerType triggerType = trigger.getType();
TenantId tenantId = triggerType.isTenantLevel() ? trigger.getTenantId() : TenantId.SYS_TENANT_ID; TenantId tenantId = triggerType.isTenantLevel() ? trigger.getTenantId() : TenantId.SYS_TENANT_ID;
notificationExecutor.submit(() -> {
try { try {
List<NotificationRule> enabledRules = notificationRulesCache.getEnabled(tenantId, triggerType); List<NotificationRule> enabledRules = notificationRulesCache.getEnabled(tenantId, triggerType);
if (enabledRules.isEmpty()) { if (enabledRules.isEmpty()) {
return; return;
} }
if (trigger.deduplicate()) { if (trigger.deduplicate()) {
enabledRules = new ArrayList<>(enabledRules); enabledRules = new ArrayList<>(enabledRules);
enabledRules.removeIf(rule -> deduplicationService.alreadyProcessed(trigger, rule)); enabledRules.removeIf(rule -> deduplicationService.alreadyProcessed(trigger, rule));
} }
final List<NotificationRule> rules = enabledRules; final List<NotificationRule> rules = enabledRules;
notificationExecutor.submit(() -> {
for (NotificationRule rule : rules) { for (NotificationRule rule : rules) {
try { try {
processNotificationRule(rule, trigger); processNotificationRule(rule, trigger);
@ -96,10 +95,10 @@ public class DefaultNotificationRuleProcessor implements NotificationRuleProcess
log.error("Failed to process notification rule {} for trigger type {} with trigger object {}", rule.getId(), rule.getTriggerType(), trigger, e); log.error("Failed to process notification rule {} for trigger type {} with trigger object {}", rule.getId(), rule.getTriggerType(), trigger, e);
} }
} }
}); } catch (Throwable e) {
} catch (Throwable e) { log.error("Failed to process notification rules for trigger: {}", trigger, e);
log.error("Failed to process notification rules for trigger: {}", trigger, e); }
} });
} }
private void processNotificationRule(NotificationRule rule, NotificationRuleTrigger trigger) { private void processNotificationRule(NotificationRule rule, NotificationRuleTrigger trigger) {

Loading…
Cancel
Save