Какие? Разве сейчас не используются все распределенные транзакции? вы все еще не знаете?

задняя часть Архитектура

Что такое распределенная транзакция

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

image.png

С помощью приведенного выше рисунка, если мы сами реализуем распределенную транзакцию, как ее реализовать?

  1. Распределенные транзакции через компенсацию
  2. Контролируется глобальными транзакциями
  3. Надежные события на основе очередей сообщений
  4. ...

Существует много способов его реализации, но как это реализовать, все равно необходимо определить в соответствии с бизнес-сценарием компании. Ниже представлены два решения сильной согласованности.

Решение с сильной консистенцией

Распределенная транзакция XA

Самая ранняя модель распределенных транзакций — это модель X/Open Distributed Transaction Processing (DTP), предложенная X/Open International Alliance, которую часто называют протоколом X/Open XA, именуемым протоколом XA. Наиболее типичной реализацией XA является протокол двухфазной фиксации (2PC).

протокол двухфазной фиксации

Как следует из названия, он описывает весь процесс в два этапа.

image.png

Поток выполнения 2PC

Первый этап:

  1. Координатор спрашивает каждого участника, можно ли нормально выполнить транзакционную операцию, и начинает ждать ответа каждого участника
  2. Каждый участник выполняет локальные транзакции (пишет локальные журналы отмены/повторения), но не фиксирует транзакции.
  3. Каждый участник возвращает ответ на запрос транзакции координатору.

Второй этап:

Отзывы всех участников Да (зафиксировать транзакцию)

  1. Координатор отправляет запрос фиксации всем участникам
  2. После того, как участник получит запрос на фиксацию, он формально выполнит операцию фиксации транзакции и отпустит ее после завершения.
  3. После того, как участник отправит сообщение подтверждения, координатор получает и завершает транзакцию.

Один или несколько участников сообщили Нет (прерванная транзакция)

  1. Координатор выдает всем участникам запрос на откат
  2. После того, как участник получит запрос на откат, он будет использовать журнал отмены, записанный на этапе 1, для выполнения операции отката транзакции и освобождения ресурсов, занятых во время выполнения всей транзакции после завершения отката.
  3. Участник отправляет сообщение подтверждения координатору после завершения отката транзакции.
  4. После того, как координатор получает информацию о подтверждении, отправленную всеми участниками, транзакция прерывается.

О первой и второй фазах 2PC сообщают все участники и координаторы, которые выполняются последовательно.

Таким образом, вы также можете найти преимущества и недостатки 2PC.

Преимущество: простой принцип

недостаток:

  1. Синхронная блокировка: во время выполнения фаз 1 и 2 транзакционные операции всех участников находятся в состоянии блокировки.
  2. Единая точка отказа: координатор — это единая точка, при возникновении проблемы с координатором весь процесс блокируется и не может быть выполнен, а механизм тайм-аута для участников и координаторов отсутствует.
  3. Несогласованность данных: после того, как координатор отправляет участникам запрос на фиксацию, если возникает дрожание сети, некоторые участники получают запрос на фиксацию, а некоторые — нет, и возникает несогласованность данных.

Из-за этих недостатков 2PC появились все протоколы трехфазной фиксации, чтобы компенсировать некоторые недостатки протокола двухфазной фиксации.

Протокол трехэтапной фиксации

Трехфазный протокол фиксации (3PC) — это улучшенная версия 2PC, которая делит процесс «запроса фиксации транзакции» 2PC на две части и становится протоколом обработки транзакций, состоящим из CanCommit, PreCommit и doCommit. И ввел механизм тайм-аута.

image.png

Процесс выполнения 3PC

Первая фаза (фаза canCommit):

  1. Координатор последовательно отправляет запросы CanCommit всем участникам. Спросите, можно ли выполнить операцию фиксации транзакции. Затем начните ждать ответа участника.
  2. После того, как участник получит запрос CanCommit, при нормальных обстоятельствах, если он считает, что транзакция может быть успешно выполнена, он вернет ответ «Да» и войдет в состояние подготовки, в противном случае он вернет «Нет».

Второй этап (этап preCommit):

Отзывы всех участников Да

  1. Координатор отправляет участнику запрос PreCommit и переходит в фазу «Подготовлено».
  2. Участник получает запрос PreCommit, выполняет операцию транзакции и записывает информацию об отмене и повторении в журнале транзакций.
  3. Если участник успешно выполнит транзакцию, он вернет ответ подтверждения и начнет ждать финальной команды.

Некоторые участники сообщили «Нет» или дождались тайм-аута.

  1. Координатор отправляет всем участникам запрос на прерывание.
  2. После того, как участник получает запрос на прерывание от координатора (или по истечении таймаута, запрос от координатора не получен), исполнение транзакции прерывается

Третий этап (этап выполнения):

выполнение успешно

  1. Когда координатор получает ответ подтверждения, отправленный участником, он переходит из состояния предварительной фиксации в состояние фиксации и отправляет запрос на выполнение всем участникам.
  2. После того, как участник получает запрос doCommit, выполняется формальная фиксация транзакции. и освободить ресурсы транзакции после завершения.
  3. После освобождения ресурсов транзакции отправьте ответ подтверждения координатору.
  4. Координатор завершает транзакцию после получения ответов подтверждения от всех участников.

прервать транзакцию

  1. Координатор отправляет всем участникам запрос на прерывание.
  2. После того, как участник получает запрос на прерывание, журнал отмены используется для выполнения операции отката транзакции, и ресурсы всех транзакций освобождаются после завершения отката.
  3. После того, как участник завершает откат транзакции, он отправляет сообщение подтверждения координатору
  4. После того, как координатор получит всю информацию о подтверждении, выполнение транзакции прерывается.

Хотя трехфазная фиксация решает проблему синхронной блокировки и отсутствия механизма тайм-аута двухфазной фиксации. Однако в трехфазном коммите остается еще много нерешенных проблем:

  1. единая точка отказа
  2. Сбои в сети между координаторами и участниками в конечном итоге приведут к несогласованности данных.

Однако по сравнению с интернет-проектами 2PC и 3PC неприменимы в случае высокой параллелизма.

Решение TCC для распределенных транзакций

TCC (Try-Confirm-Cancel) также известна как двухэтапная компенсационная транзакция, и TCC имеет большое количество приложений в Ant Financial.

  • Этап 1. Попробуйте: Обнаружение зарезервированных ресурсов
  • Этап 2 - Подтверждение: представлена ​​реальная бизнес-операция
  • Фаза 3 — Отмена: освобождение зарезервированных ресурсов

image.pngВсе эти три бизнес-логики должны быть реализованы в бизнес-логике.

image.png

это процесс заказа

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

image.pngЗаказ успешно размещен, и измененный инвентарь равен 98, но вычитаемый инвентарь должен быть изменен на 2 в замороженном инвентаре, и то же самое верно для других. Выполните локальную транзакцию после завершения. Это то, что делает фаза Try. Если все успешно, перейдите к этапу подтверждения

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

image.pngОднако если транзакция откатывается на этапе попытки, это этап отмены.

Стадия отмены: в это время структура транзакции TCC будет воспринимать это. Поэтому все сервисы будут уведомлены об откате, то есть о пополнении замороженного инвентаря обратно в инвентарь, и то же самое верно и для других. На этапе отмены, если транзакционная операция не удалась, ее необходимо все время повторять.

image.png

Некоторые оптимизации TCC объясняются в отдельной главе.

Самостоятельно разработанная распределенная транзакция: локальная таблица транзакций + решение для окончательной согласованности сообщений

Сценарии с сильной согласованностью, такие как 2PC, на самом деле очень узкие.В большинстве сценариев необходимо только убедиться, что данные в конечном итоге непротиворечивы. После небольшой задержки данные в конечном итоге согласуются. Например, в процессе заказа, после того, как заказ был успешно размещен, нам нужно очистить корзину и увеличить количество баллов, На самом деле это нормально на несколько секунд позже. Нам нужно только убедиться, что заказ выполнен успешно, а корзина покупок очищена или купоны, красные конверты и накопленный опыт — все успешно или неудачно.

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

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

image.png

Как упоминалось выше, мы будем реализовывать его шаг за шагом.

Сначала вам нужно иметь класс сущности для локальных сообщений:

/**
 * 事务消息实体  和数据库是一一对应的字段
 */
public class MsgInfo {
    /**
     * 主键
     */
    private Long id;
    
    /**
     * 事务消息
     */
    private String content;
    
    /**
     * 主题和RocketMQ对应
     */
    private String topic;
    
    /**
     * 标签和RocketMQ对应
     */
    private String tag;
    
    /**
     * 状态:1-等待,2-发送
     */
    private int status;
    
    /**
     * 创建时间
     */
    private Date createTime;
    
    /**
     * 延迟时间(单位:s)
     * 最迟几秒钟要发的MQ中
     */
    private int delay;
    
 }

После создания класса сообщения сущности здесь необходим класс для работы с сущностью.

/**
 * 事务操作实体 实际上就是一个Queue
 */
public class Msg {
    /**
     * 主键 和MsgInfo是一一对应的
     */
    private Long id;
    
    /**
     * db-url key, 跟 数据源map做映射
     */
    private String url;
    
    /**
     * 已经处理次数
     */
    private int haveDealedTimes;
    
    /**
     * 创建时间
     */
    private long createTime;
    
    /**
     * 下次超时时间
     */
    private long nextExpireTime;
    
 }

После того, как класс сущности создан, давайте подумаем.Согласно шагам, которые мы сказали ранее, нам нужно поместить заказ в MQ, но здесь мы вообще не отправляем его напрямую в MQ.После того, как мы напишем БД, мы сначала нужно отправить его в очередь.В очереди класс сущности Msg хранится в очереди, поэтому здесь нам нужен рабочий поток доставки, чтобы всплывать данные из очереди, и рабочий поток доставки получает их и сравнивает это с таблицей сообщений в базе данных, чтобы увидеть, есть ли в базе данных какие-либо локальные сообщения о транзакциях. Если нет, сообщение заканчивается. Если есть, то нужно создать MQ сообщение, перед созданием сообщения нужно установить максимальное количество повторов для обеспечения точности (может случиться так, что транзакция не была отправлена ​​на чтение или считывание), может повторить Есть шанс попробовать, если лимит не превышен и запрос получен, то для доставки нужно создать MQ сообщение. Завершить немедленно в случае успеха.

image.png

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

image.png

Если в это время сервис завис и все в очереди пропало, что нам делать в это время? Так что на этот раз нам также понадобится нить-ловушка. После перезапуска поток перехвата проверяет, удерживается ли блокировка.Если блокировка удерживается, поток ожидающих транзакций за последние десять минут из БД будет оцениваться в соответствии со статусом MsgInfo, а затем будет создано сообщение MQ. и доставлено. Если доставка прошла успешно, она завершится. Если доставка не удалась, она снова попадет в очередь повторных попыток.

image.png

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

image.png

Давайте посмотрим на код инициализации:

/**
 * @方法名称 init
 * @功能描述 init 初始化,config才是ok
 * @param config 配置对象
 */
public void init(Config config) {
    if (state.get().equals(State.RUNNING)) {
        LOGGER.info("Msg Processor have inited return");
        return;
    }
    LOGGER.info("MsgProcessor init start");
    state.compareAndSet(State.CREATE, State.RUNNING);
    // 1、设置环境
    this.serverType = ConfigUtil.getServerType();
    if (config.getEtcdHosts() != null && config.getEtcdHosts().length >= 1) {
        envNeedLock = true;
        defaultEtcd = config.getEtcdHosts();
        LOGGER.info("serverType {} envNeedLock {} etcdhosts {}", serverType, envNeedLock, defaultEtcd);
    }
    // 2、设置配置
    this.config = config;
    // 3、设置 事务消息处理线程数
    exeService = Executors.newFixedThreadPool(config.getThreadNum(), new ThreadFactory("MsgProcessorThread-"));
    for (int i = 0; i < config.getThreadNum(); i++) {
        exeService.submit(new MsgDeliverTask());
    }
    // 4、设置 其他线程
    scheService = Executors.newScheduledThreadPool(config.getSchedThreadNum(), new ThreadFactory("MsgScheduledThread-"));
    // 设置时间转动线程:时间轮重试投递失败的事务操作
    scheService.scheduleAtFixedRate(new TimeWheelTask(), TIME_WHEEL_PERIOD, TIME_WHEEL_PERIOD, TimeUnit.MILLISECONDS);
    // 设置事务消息删除线程
    scheService.scheduleAtFixedRate(new CleanMsgTask(), config.deleteTimePeriod, config.deleteTimePeriod, TimeUnit.SECONDS);
    // 设置 补漏线程:防止最近10分钟的线程被漏提交
    scheService.scheduleAtFixedRate(new ScanMsgTask(), config.schedScanTimePeriod, config.schedScanTimePeriod, TimeUnit.SECONDS);
    // 设置心跳线程:汇报 事务提交队列的堆积情况
    scheService.scheduleAtFixedRate(new Runnable() {
        @Override
        public void run() {
            LOGGER.info("stats info msgQueue size {} timeWheelQueue size {}", msgQueue.size(), timeWheel.size());
        }
    }, 20, config.getStatsTimePeriod(), TimeUnit.SECONDS);
    // 6、初始化锁客户端
    initLock();
    LOGGER.info("MsgProcessor init end");
}

Ветка доставки транзакции:

//事务投递线程
class MsgDeliverTask implements Runnable {
    @Override
    public void run() {
        while (true) {
            if (!state.get().equals(State.RUNNING)) {
                break;
            }
            try {
                // 1、每100ms从 队列 弹出一条事务操作消息
                Msg msg = null;
                try {
                    //拉出来消息
                    msg = msgQueue.poll(DEF_TIMEOUT_MS, TimeUnit.MILLISECONDS);
                } catch (InterruptedException ex) {
                }
                if (msg == null) {
                    continue;
                }
                LOGGER.debug("poll msg {}", msg);
                int dealedTime = msg.getHaveDealedTimes() + 1;
                msg.setHaveDealedTimes(dealedTime);
                // 2、从db获取实际事务消息(这里我们不知道是否事务已经提交,所以需要从DB里面拿)
                MsgInfo msgInfo = msgStorage.getMsgById(msg);
                LOGGER.debug("getMsgInfo from DB {}", msgInfo);
                if (msgInfo == null) {
                    if (dealedTime < MAX_DEAL_TIME) {
                        // 3.1、加入时间轮转动队列:重试投递
                        long nextExpireTime = System.currentTimeMillis() + TIMEOUT_DATA[dealedTime];
                        msg.setNextExpireTime(nextExpireTime);
                        timeWheel.put(msg);
                        LOGGER.debug("put msg in timeWhellQueue {} ", msg);
                    }
                } else {
                    // 3.2、投递事务消息
                    Message mqMsg = buildMsg(msgInfo);
                    LOGGER.debug("will sendMsg {}", mqMsg);
                    SendResult result = producer.send(mqMsg);
                    LOGGER.info("msgId {} topic {} tag {} sendMsg result {}", msgInfo.getId(), mqMsg.getTopic(), mqMsg.getTags(), result);
                    if (null == result || result.getSendStatus() != SendStatus.SEND_OK) {
                        // 投递失败,重入时间轮
                        if (dealedTime < MAX_DEAL_TIME) {
                            long nextExpireTime = System.currentTimeMillis() + TIMEOUT_DATA[dealedTime];
                            msg.setNextExpireTime(nextExpireTime);
                            timeWheel.put(msg);
                            // 这里可以优化 ,因为已经确认事务提交了,可以从DB中拿到了
                            LOGGER.debug("put msg in timeWhellQueue {} ", msg);
                        }
                    } else if (result.getSendStatus() == SendStatus.SEND_OK) {
                        // 投递成功,修改数据库的状态(标识已提交)
                        int res = msgStorage.updateSendMsg(msg);
                        LOGGER.debug("msgId {} updateMsgStatus success res {}", msgInfo.getId(), res);
                    }
                }
            } catch (Throwable t) {
                LOGGER.error("MsgProcessor deal msg fail", t);
            }
        }
    }
}

Поток колеса времени повтора

class TimeWheelTask implements Runnable {
    @Override
    public void run() {
        try {
            if (state.get().equals(State.RUNNING)) {
                long cruTime = System.currentTimeMillis();
                //检查是否有Msg
                Msg msg = timeWheel.peek();
                // 拿出来的时候有可能还没有超时
                while (msg != null && msg.getNextExpireTime() <= cruTime) {
                    msg = timeWheel.poll();
                    LOGGER.debug("timeWheel poll msg ,return to msgQueue {}", msg);
                    // 重新放进去
                    msgQueue.put(msg);
                    msg = timeWheel.peek();
                }
            }
        } catch (Exception ex) {
            LOGGER.error("pool timequeue error", ex);
        }
    }
}

удалить локальную ветку сообщений

class CleanMsgTask implements Runnable {
    @Override
    public void run() {
        if (state.get().equals(State.RUNNING)) {
            LOGGER.debug("DeleteMsg start run");
            try {
                Iterator<DataSource> it = msgStorage.getDataSourcesMap().values().iterator();
                while (it.hasNext()) {
                    DataSource dataSrc = it.next();
                    if (holdLock) {
                        LOGGER.info("DeleteMsgRunnable run ");
                        int count = 0;
                        int num = config.deleteMsgOneTimeNum;
                        // 影响行数 不等于 删除数 及 大于最大删除数时,本次task结束
                        while (num == config.deleteMsgOneTimeNum && count < MAX_DEAL_NUM_ONE_TIME) {
                            try {
                                num = msgStorage.deleteSendedMsg(dataSrc, config.deleteMsgOneTimeNum);
                                count += num;
                            } catch (SQLException e) {
                                LOGGER.error("deleteSendedMsg fail ", e);
                            }
                        }
                    }
                }
            } catch (Exception ex) {
                LOGGER.error("delete Run error ", ex);
            }
        }
    }
}

Ловушка нить

class ScanMsgTask implements Runnable {
    @Override
    public void run() {
        if (state.get().equals(State.RUNNING)) {
            LOGGER.debug("SchedScanMsg start run");
            Iterator<DataSource> it = msgStorage.getDataSourcesMap().values().iterator();
            while (it.hasNext()) {
                DataSource dataSrc = it.next();
                boolean canExe = holdLock;
                if (canExe) {
                    LOGGER.info("SchedScanMsgRunnable run");
                    int num = LIMIT_NUM;
                    int count = 0;
                    while (num == LIMIT_NUM && count < MAX_DEAL_NUM_ONE_TIME) {
                        try {
                            List<MsgInfo> list = msgStorage.getWaitingMsg(dataSrc, LIMIT_NUM);
                            num = list.size();
                            if (num > 0) {
                                LOGGER.debug("scan db get msg size {} ", num);
                            }
                            count += num;
                            for (MsgInfo msgInfo : list) {
                                try {
                                    Message mqMsg = buildMsg(msgInfo);
                                    SendResult result = producer.send(mqMsg);
                                    LOGGER.info("msgId {} topic {} tag {} sendMsg result {}", msgInfo.getId(), mqMsg.getTopic(), mqMsg.getTags(), result);
                                    if (result != null && result.getSendStatus() == SendStatus.SEND_OK) {
                                        // 修改数据库的状态
                                        int res = msgStorage.updateMsgStatus(dataSrc, msgInfo.getId());
                                        LOGGER.debug("msgId {} updateMsgStatus success res {}", msgInfo.getId(), res);
                                    }
                                } catch (Exception e) {
                                    LOGGER.error("SchedScanMsg deal fail", e);
                                }
                            }
                        } catch (SQLException e) {
                            LOGGER.error("getWaitMsg fail", e);
                        }
                    }
                }
            }
        }
    }
    
}

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