如何通过编程检测AEM CaaS中复制队列是否阻塞
AEM CaaS复制队列阻塞邮件通知解决方案
先分析你现有两种方法的问题:
- 方法1:
getNextRetryTime() !=0的判断逻辑太宽泛,只要队列存在待重试任务就触发通知,但临时网络波动导致的正常重试不属于真正的队列阻塞;再加上ReplicationEventHandler会在每次复制事件(新增任务、重试触发)时执行,所以会频繁重复发送邮件。 - 方法2:
firstEntry.getNumProcessed()统计的是任务被处理的总次数,和失败重试次数不直接对应;而且代码里entries().isEmpty()的判断逻辑写反了(else块注释是"Queue is Not Empty",实际else对应队列空的场景),导致真正阻塞时无法触发判断。
可行解决方案核心思路
要实现准确的阻塞通知,需要做到两点:精准判断队列阻塞状态+避免重复发送邮件,以下是具体实现:
核心判断逻辑
同时覆盖两种典型阻塞场景:
- 复制队列被系统标记为
BLOCKED状态(比如代理连接失败被暂停) - 队列中有任务多次重试失败(达到设定阈值)且处于失败状态,导致队列堆积
代码实现示例
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发送模板邮件 } }
关键说明
- 事件过滤:只在复制失败或队列状态变化时执行检查,减少不必要的系统开销
- 双维度阻塞判断:同时覆盖队列主动标记阻塞和任务重试失败堆积两种场景,避免漏判
- 去重机制:通过JCR节点存储上次发送时间,确保设定间隔内只发一次,避免邮件轰炸
- CaaS环境适配:使用
ServiceResourceResolver访问JCR存储,符合CaaS环境的权限管控要求,需提前配置对应的subservice权限
额外注意事项
- 确保OSGi组件正确注册,依赖的
AgentManager、ResourceResolverFactory等服务正常注入 - 提前配置AEM的
MailService,确保SMTP参数正确 - 可根据实际业务需求调整
EMAIL_INTERVAL_MINUTES和RETRY_FAILURE_THRESHOLD的值
内容的提问来源于stack exchange,提问作者Vivek Dhiman
相关产品推荐
相关产品推荐

