You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过编程检测AEM CaaS中复制队列是否阻塞

AEM CaaS复制队列阻塞邮件通知解决方案

先分析你现有两种方法的问题:

  • 方法1:getNextRetryTime() !=0 的判断逻辑太宽泛,只要队列存在待重试任务就触发通知,但临时网络波动导致的正常重试不属于真正的队列阻塞;再加上ReplicationEventHandler会在每次复制事件(新增任务、重试触发)时执行,所以会频繁重复发送邮件。
  • 方法2:firstEntry.getNumProcessed() 统计的是任务被处理的总次数,和失败重试次数不直接对应;而且代码里entries().isEmpty()的判断逻辑写反了(else块注释是"Queue is Not Empty",实际else对应队列空的场景),导致真正阻塞时无法触发判断。

可行解决方案核心思路

要实现准确的阻塞通知,需要做到两点:精准判断队列阻塞状态+避免重复发送邮件,以下是具体实现:

核心判断逻辑

同时覆盖两种典型阻塞场景:

  1. 复制队列被系统标记为BLOCKED状态(比如代理连接失败被暂停)
  2. 队列中有任务多次重试失败(达到设定阈值)且处于失败状态,导致队列堆积

代码实现示例

import org.apache.sling.api.resource.ResourceResolver;
import org.apache.sling.api.resource.ResourceResolverFactory;
import com.day.cq.replication.*;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.TimeUnit;
import org.apache.sling.api.resource.ValueMap;

public class ReplicationBlockedNotificationHandler implements ReplicationEventHandler {
    private static final Logger log = LoggerFactory.getLogger(ReplicationBlockedNotificationHandler.class);
    // 邮件发送间隔(避免轰炸),单位:分钟
    private static final long EMAIL_INTERVAL_MINUTES = 30;
    // 任务重试失败阈值
    private static final int RETRY_FAILURE_THRESHOLD = 5;
    // JCR存储上次发送通知的时间戳路径(CaaS环境适配)
    private static final String LAST_NOTIFICATION_PATH = "/var/replication-notifications/publish-blocked";

    private final ResourceResolverFactory resourceResolverFactory;

    // OSGi构造注入资源解析工厂
    public ReplicationBlockedNotificationHandler(ResourceResolverFactory resourceResolverFactory) {
        this.resourceResolverFactory = resourceResolverFactory;
    }

    @Override
    public void handleEvent(ReplicationEvent event) {
        // 只在复制失败或队列状态变化时触发检查,减少无效执行
        if (!event.getType().equals(ReplicationEvent.Type.REPLICATION_FAILED) && 
            !event.getType().equals(ReplicationEvent.Type.REPLICATION_QUEUE_CHANGED)) {
            return;
        }

        AgentManager agentManager = ReplicationUtils.getAgentManager(event.getResourceResolver());
        if (agentManager == null) {
            log.error("无法获取Agent Manager");
            return;
        }

        Agent publishAgent = agentManager.getAgent("publish");
        // 跳过无效/未启用的代理
        if (publishAgent == null || !publishAgent.isEnabled() || !publishAgent.isValid()) {
            return;
        }

        ReplicationQueue replicationQueue = publishAgent.getQueue();
        if (replicationQueue == null) {
            return;
        }

        // 判断队列是否真正阻塞
        boolean isQueueBlocked = false;
        // 场景1:队列状态直接标记为阻塞
        if (ReplicationQueue.State.BLOCKED.equals(replicationQueue.getState())) {
            isQueueBlocked = true;
        } 
        // 场景2:队列有任务,且首个任务多次重试失败
        else if (!replicationQueue.entries().isEmpty()) {
            ReplicationQueue.Entry firstEntry = replicationQueue.entries().get(0);
            if (ReplicationQueue.Entry.Status.FAILED.equals(firstEntry.getStatus()) && 
                firstEntry.getNumRetries() > RETRY_FAILURE_THRESHOLD) {
                isQueueBlocked = true;
            }
        }

        if (isQueueBlocked) {
            // 检查是否满足发送间隔要求
            if (shouldSendNotification()) {
                Map<String, String> emailParams = new HashMap<>();
                emailParams.put("agentId", publishAgent.getId());
                emailParams.put("agent配置路径", publishAgent.getConfiguration().getConfigPath());
                emailParams.put("队列状态", replicationQueue.getState().name());
                emailParams.put("待处理任务数", String.valueOf(replicationQueue.entries().size()));
                
                sendEmail(emailParams);
                updateLastNotificationTime();
                log.info("::: 复制队列已阻塞 - 已发送通知邮件 :::");
            } else {
                log.info("::: 复制队列已阻塞 - 通知邮件已在间隔内发送,无需重复发送 :::");
            }
        }
    }

    // 检查是否符合邮件发送间隔要求
    private boolean shouldSendNotification() {
        try (ResourceResolver resolver = resourceResolverFactory.getServiceResourceResolver(
                Map.of(ResourceResolverFactory.SUBSERVICE, "replication-notification"))) {
            if (resolver.resourceExists(LAST_NOTIFICATION_PATH)) {
                long lastSentTime = resolver.getResource(LAST_NOTIFICATION_PATH)
                                            .adaptTo(ValueMap.class)
                                            .get("lastSent", 0L);
                long currentTime = System.currentTimeMillis();
                return currentTime - lastSentTime > TimeUnit.MINUTES.toMillis(EMAIL_INTERVAL_MINUTES);
            }
            // 首次发送直接允许
            return true;
        } catch (Exception e) {
            log.error("检查上次通知时间失败", e);
            // 出错时默认允许发送,避免遗漏通知
            return true;
        }
    }

    // 更新上次发送通知的时间戳到JCR
    private void updateLastNotificationTime() {
        try (ResourceResolver resolver = resourceResolverFactory.getServiceResourceResolver(
                Map.of(ResourceResolverFactory.SUBSERVICE, "replication-notification"))) {
            if (!resolver.resourceExists(LAST_NOTIFICATION_PATH)) {
                // 创建存储节点
                resolver.create(resolver.getResource("/var/replication-notifications"), 
                                "publish-blocked", Map.of());
            }
            ValueMap valueMap = resolver.getResource(LAST_NOTIFICATION_PATH).adaptTo(ValueMap.class);
            valueMap.put("lastSent", System.currentTimeMillis());
            resolver.commit();
        } catch (Exception e) {
            log.error("更新上次通知时间失败", e);
        }
    }

    // 邮件发送逻辑(需结合AEM MailService实现)
    private void sendEmail(Map<String, String> params) {
        // 此处实现具体邮件发送,比如调用AEM的MailService发送模板邮件
    }
}

关键说明

  1. 事件过滤:只在复制失败或队列状态变化时执行检查,减少不必要的系统开销
  2. 双维度阻塞判断:同时覆盖队列主动标记阻塞和任务重试失败堆积两种场景,避免漏判
  3. 去重机制:通过JCR节点存储上次发送时间,确保设定间隔内只发一次,避免邮件轰炸
  4. CaaS环境适配:使用ServiceResourceResolver访问JCR存储,符合CaaS环境的权限管控要求,需提前配置对应的subservice权限

额外注意事项

  • 确保OSGi组件正确注册,依赖的AgentManager、ResourceResolverFactory等服务正常注入
  • 提前配置AEM的MailService,确保SMTP参数正确
  • 可根据实际业务需求调整EMAIL_INTERVAL_MINUTES和RETRY_FAILURE_THRESHOLD的值

内容的提问来源于stack exchange,提问作者Vivek Dhiman

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.14 22:30:33