Серия DelayQueue (3): решение для сохраняемости

задняя часть

Первоначально опубликовано в ЦзяньшуСхема сохраняемости DelayQueue, Это обновление в основном оптимизирует метод processTask и оптимизирует настройку дополнительного времени задержки для выполнения компенсации, которую можно прочитать для сравнения.

В предыдущей статье мы упоминали, что использовали DelayQueue в проекте для решения некоторых задач, которые необходимо отложить, но недавно столкнулись с проблемой в производственной среде. После перезапуска сервера эти невыполненные отложенные задачи исчезают. Итак, на повестке дня стоит вопрос о том, как сохранить отложенные задачи.

Что касается конкретной схемы реализации DelayQueue, то она была в предыдущей статье.Серия DelayQueue (2): основные компонентыупоминается в. Эта статья не будет повторяться.

Темой этого выпуска является обсуждение того, как сохранить отложенные задачи.

Какова устойчивость отложенных задач? Как следует из названия, он предназначен для хранения необходимых данных для выполнения этих отложенных задач в базе данных или Redis.

Так зачем упорствовать? Это очень просто, т.к. данные отложенной задачи хранятся в памяти, то нам нужно прописать свои постоянные записи для достижения высокой доступности, иначе при сбое сервера или только что выпущенная версия вызывает перезагрузку сервера, то неисполненные Данные отложенной миссии будут полностью потеряны, чего мы явно не хотим.

Мой текущий план таков:
1. Если необходимо использовать DelayQueue, вызовите метод saveDelayTask.Обязательные параметры включают тег маршрутизации класса фабрики политики функции задержки задачи, параметр messageBody в формате json, требуемый методом выполнения, и время задержки выполнения. задержкиВремя в секундах.
2. Планировщик задач выполняет метод getNotCompletedMessageList каждые 5 секунд.

В большинстве случаев метод processTask будет выполняться вовремя в ожидаемое время выполнения, тогда в нештатных условиях, если сервер перезагрузится, планирование запланированных задач найдет те отложенные задачи, которые не выполняются по расписанию через определенный период времени. , Метод по очереди выполняет методы processTask соответствующих задач.

В ненормальном состоянии выполнение отложенной задачи будет отложено до определенной степени, чем ожидаемое время выполнения.Разработанный мной план находится в пределах диапазона, который мы можем разрешить в настоящее время.Вы можете установить время задержки альтернативного плана как соответствующий.

В течение этого периода будет использоваться таблица yb_delay_task_message.Следующая структура таблицы:

CREATE TABLE `yb_delay_task_message` (
  `id` bigint(20) unsigned NOT NULL AUTO_INCREMENT COMMENT '自增id',
  `tag` varchar(128) NOT NULL DEFAULT '' COMMENT '延迟队列执行函数的key',
  `message_body` longtext CHARACTER SET utf8mb4 COMMENT '消息体,以json格式存储',
  `status` tinyint(3) unsigned DEFAULT '0' COMMENT '状态;0:未完成,1:已完成,2:已失败 3:执行中',
  `error_stack` longtext COMMENT '失败堆栈',
  `version` bigint(20) unsigned NOT NULL DEFAULT '0' COMMENT '版本号',
  `ip_address` bigint(20) DEFAULT NULL COMMENT '执行ip地址',
  `delay_time` bigint(20) NOT NULL COMMENT '延迟执行的时间长度',
  `expected_time` datetime DEFAULT NULL COMMENT '预计执行时间',
  `execution_time` datetime DEFAULT NULL COMMENT '实际执行时间',
  `create_time` datetime NOT NULL COMMENT '创建时间',
  `update_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '修改时间',
  PRIMARY KEY (`id`),
  KEY `idx_expectedTime_status` (`expected_time`,`status`)
) ENGINE=InnoDB  DEFAULT CHARSET=utf8 COMMENT='延迟队列消息表';

Основной код примерно такой, остальные коды очень простые и не будут публиковаться один за другим.

public void saveDelayTask(String tag, String messageBody, Long delayTime) {
    DelayTaskMessage delayTaskMessage = new DelayTaskMessage();
    delayTaskMessage.setTag(tag);
    LocalDateTime now = LocalDateTime.now();
    delayTaskMessage.setCreateTime(now);
    delayTaskMessage.setUpdateTime(now);
    delayTaskMessage.setDelayTime(delayTime);
    delayTaskMessage.setExpectedTime(now.plusSeconds(delayTime));
    delayTaskMessage.setMessageBody(messageBody);
    delayTaskMessage.setStatus(KafkaMessageStatusEnum.NOT_COMPLETE.getCode());
    int res = delayTaskMessageMapper.insertDelayTaskMessage(delayTaskMessage);
    if (res <= 0) {
        log.error("ybBrokerApp|insertDelayTaskMessage error, res<=0");
        throw new RuntimeException("insertDelayTaskMessage error, res<=0");
    }
    TaskMessage taskMessage = new TaskMessage(delayTime * 1000, messageBody,
            function -> this.processTask(delayTaskMessage));
    DelayQueue<TaskMessage> queue = taskManager.getQueue();
    queue.offer(taskMessage);
}

Во-первых, давайте проанализируем метод saveDelayTask, используемый для сохранения отложенных задач.

тег относится к тегу отложенной задачи, который используется для указания соответствующего класса политики.
messageBody в основном используется для хранения некоторых необходимых данных для выполнения отложенных задач и хранится в методе json.
delayTime — это время задержки, единица измерения по умолчанию — с, в основном для простоты использования.

Основная функция этого метода состоит в том, чтобы сначала сохранить отложенную задачу, которая еще не была выполнена, автоматически рассчитать ожидаемое время выполнения отложенной задачи в соответствии со временем задержки, чтобы облегчить отслеживание последующего алгоритма компенсации, а затем использовать характеристики DelayQueue для отправки отложенной задачи в очередь задержки.

public int processTask(DelayTaskMessage param) {
    DelayTaskMessage delayTaskMessage = delayTaskMessageMapper.getDelayTaskMessageById(param.getId());
    try {
        if (null != delayTaskMessage) {
            if (!Objects.equals(delayTaskMessage.getStatus(), DelayTaskMessageStatusEnum.NOT_COMPLETE.getCode())) {
                log.info("processTask executed already");
                return 1;
            }
            try {
                delayTaskMessage.setIpAddress(InetAddress.getLocalHost().getHostAddress());
            } catch (UnknownHostException ex) {
                log.error("Address.getLocalHost error", ex);
            }
            int res = delayTaskMessageMapper.delayTaskStartProcess(delayTaskMessage);
            if (res <= 0) {
                log.info("delayTaskStartProcess error,maybe processTask executed already");
                return 1;
            }
            //处理逻辑
            DelayTaskExecuteProcessor delayTaskExecuteProcessor = delayTaskExecuteProcessorFactory.getExecuteProcessor(delayTaskMessage.getTag());
            if (delayTaskExecuteProcessor != null) {
                delayTaskExecuteProcessor.execute(delayTaskMessage);
            } else {
                throw new RuntimeException("no such processor,tag=" + delayTaskMessage.getTag());
            }
            delayTaskMessage.setExecutionTime(LocalDateTime.now());
            res = delayTaskMessageMapper.delayTaskProcessSuccess(delayTaskMessage);
            if (res <= 0) {
                log.error("delayTaskProcessSuccess error");
                return 1;
            }
            return 1;
        } else {
            log.error("ybBrokerApp processTask error, delayTaskMessage is null delayTaskMessageId=", param.getId());
            return 0;
        }
    } catch (Exception e) {
        log.error("ybBrokerApp processTask error , param = " + param.toString() + "|", e);
        if (null != delayTaskMessage) {
            delayTaskMessage.setErrorStack(e.getMessage());
            delayTaskMessageMapper.delayTaskProcessFail(delayTaskMessage);
        }
        return 0;
    }
}

Кроме того, есть основной метод processTask, который обрабатывает отложенные задачи.

1. По id найти в базе постоянные данные, соответствующие отложенной задаче, которую необходимо выполнить.
2. Если постоянные данные не пусты, а состояние не невыполнено, будет указано, что задачу не нужно выполнять снова, чтобы предотвратить повторное выполнение. Статус здесь имеет в общей сложности четыре состояния: не выполняется, выполняется, выполняется успешно и выполняется с ошибкой.
3. Если постоянные данные не пусты и не выполняются, то измените состояние выполнения данных на выполнение, и запишите ip-адрес метода выполнения для последующего анализа, а затем контролируйте его через проблемы Concurrency версии. Только когда версия совпадает с версией в базе данных, а статус в базе данных не выполняется, разрешается изменить статус на выполнение.
4. Если выполнение предыдущего шага прошло успешно, найдите класс стратегии, соответствующий тегу, и выполните соответствующий метод execute.
5. Затем измените статус этих постоянных данных на статус успешно выполнено.Здесь нет необходимости ограничиваться версией и статусом, и его можно напрямую изменить на успешное выполнение.
6. Если в процессе выполнения возникают другие исключения, измените статус данных на Ошибка выполнения.

public List<DelayTaskMessage> getNotCompletedMessageList(int total, int index) {
    LocalDateTime expectedTime = LocalDateTime.now().plusSeconds(1L);
    List<DelayTaskMessage> delayTaskMessageList = delayTaskMessageMapper.getNotCompletedMessageList(expectedTime,total, index);
    if (CollectionUtils.isEmpty(delayTaskMessageList)) {
        return Lists.newArrayList();
    }
    return delayTaskMessageList;
}

Последнее — это реализация плана компенсации.Я в запланированном задании гарантирую, что отложенное задание будет выполнено хотя бы один раз. Что касается того, будет ли он выполняться повторно, я контролирую это в методе processTask.

Мой план состоит в том, чтобы пройти те отложенные задачи, которые не были выполнены после ожидаемого времени выполнения + 1 с каждые 5 с. Отложенные задачи в этих списках затем вызываются обратно в метод processTask.

Если отложенная задача наконец выполняется с помощью схемы компенсации, она будет выполнена на 1–6 с позже ожидаемого времени выполнения. В настоящее время в нашем проекте эта дополнительная задержка приемлема. Вы по-прежнему должны изменить это дополнительное время задержки в соответствии с реальной ситуацией.

Я записал ожидаемое время (оценочное время выполнения) и executeTime (фактическое время выполнения) в таблицу yb_delay_task_message, чтобы по этим двум полям можно было сравнить конкретную производительность выполнения.

Выше приведена схема сохраняемости, которую я разработал для DelayQueue. Если у вас есть лучшие мнения, вы можете обсудить ее вместе.