Поговорим о небольшом понимании и решениях распределенных транзакций

задняя часть
Поговорим о небольшом понимании и решениях распределенных транзакций

помещение

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

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

Во-первых, при разделении системы почти всегда возникает проблема распределенных транзакций.Смоделированный случай выглядит следующим образом:

j-t-s-i-a-1.png
j-t-s-i-a-1.png

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

Прямые вызовы RPC в транзакциях обеспечивают высокую согласованность

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

[订单微服务请求钱包微服务进行扣款并更新订单状态]

处理订单微服务请求钱包微服务进行扣款并更新订单状态方法(){
    [开启事务]
    1、查询订单
    2、HTTP调用钱包微服务扣款
    3、更新订单状态为扣款成功
    [提交事务]
}

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

  • 1. На втором шаге вышеописанного метода, по разным причинам самого микросервиса кошелька, ответ дебетового интерфейса крайне медленный, что приведет к зависанию транзакции вышеописанного метода обработки (точнее, подключения к базе данных) в течение длительного времени, и удерживаемое соединение с базой данных не может быть освобождено, это приведет к исчерпанию соединения пула соединений с базой данных, и это может легко привести к тому, что другие зависящие от базы данных интерфейсы микрослужбы заказа перестанут отвечать.
  • 2. Микросервис кошелька представляет собой развертывание с одним узлом (не все микросервисы компании идеальны), приложение не работает во время обновления, а вызов интерфейса на шаге 2 вышеописанного метода напрямую завершается сбоем, что приведет к короткому замыканию всех транзакций. период времени.Откат, эквивалентный вычету записи заказа микросервиса, недоступен.
  • 3."сеть ненадежна", если происходит отключение сети при совершении HTTP-вызовов или получении ответов, статус между сервисами может быть непонятен друг другу, например, микросервис заказа успешно вызывает микросервис кошелька, а при получении ответа возникает проблема с сетью, и списание пройдет успешно.Нет возможности обновить статус заказа (транзакция микросервиса заказа откатывается).
j-t-s-i-a-2.png
j-t-s-i-a-2.png

Хотя сейчас естьHystrixДругие фреймворки могут изолировать вызовы на основе пулов потоков или быстро сбоить на основе автоматических выключателей, но это малоэффективно. Поэтому я лично считаю"Совершенно нежелательно добиваться строгой согласованности с прямыми вызовами RPC внутри транзакции", если этот метод используется для реализации «распределенной транзакции», рекомендуется исправить, в противном случае вы можете только молиться, чтобы нижестоящие службы или сеть не имели проблем каждый день.

Асинхронная отправка сообщения в транзакцию

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

[订单微服务请求钱包微服务进行扣款并更新订单状态]

处理订单微服务请求钱包微服务进行扣款并更新订单状态方法(){
    [开启事务]
    1、查询订单
    2、推送钱包微服务扣款消息(推送消息)
    3、更新订单状态为扣款成功
    [提交事务]
}

Если вышеуказанный метод обработки является абстрактным, он выражается следующим образом:

方法(){
    DataSource  dataSource = xx;
    Connection con = dataSource.getConnection();
    con.setAutoCommit(false);
    try{
       1、SQL操作;
       2、推送消息;
       3、SQL操作;
       con.commit();
    }catch(Exception e){
        con.rollback();
    }finally{
        释放其他资源;
        release(con);
    }
}

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

  • 1. Промежуточное ПО очереди сообщений ненормально и не может быть вызвано нормально.Общая ситуация заключается в том, что сетевая причина или промежуточное ПО очереди сообщений недоступно, что вызовет исключение и приведет к откату транзакции. Эта ситуация кажется разумной, но подумайте хорошенько: почему исключение вызова промежуточного программного обеспечения очереди сообщений приводит к откату бизнес-транзакции?Если промежуточное программное обеспечение не восстанавливается, не эквивалентен ли этот интерфейсный вызов недоступности?
  • 2. Если промежуточное ПО очереди сообщений нормальное и сообщение проталкивается нормально, но на шаге 3 транзакция откатывается из-за синтаксической ошибки в SQL, что вызовет проблему, что нижестоящий микросервис вызывается успешно, но локальный транзакция откатывается, что приводит к описанной выше проблеме: данные нижестоящей системы несовместимы.
j-t-s-i-a-3.png
j-t-s-i-a-3.png

В целом:"Асинхронная отправка сообщений в транзакцию — ненадежная реализация.".

Решения, предлагаемые в настоящее время отраслью

Текущие основные решения для распределенных транзакций в отрасли в основном включают в себя: многоэтапную схему отправки (2PC, 3PC), компенсационную транзакцию (TCC) и транзакцию сообщений (в основном RocketMQ, основная идея также заключается в многоэтапной схеме отправки и на основе промежуточного программного обеспечения). опрос и повторная попытка, другое ПО промежуточного слоя очереди сообщений не реализует распределенные транзакции). Принципы работы этих схем здесь не раскрываются, в настоящее время в сети имеется множество соответствующих материалов, резюмируем их характеристики:

  • Схема многоэтапной фиксации: для обычных двухэтапных и трехэтапных транзакций фиксации требуются дополнительные менеджеры ресурсов для координации транзакций, а согласованность данных строгая, но схема реализации более сложная, а жертва производительности относительно велика (в основном требуется блокировка ресурсы, подождите, пока все транзакции будут зафиксированы перед разблокировкой), не подходит для сценариев с высоким параллелизмом, в настоящее время более известными являются Alibaba с открытым исходным кодом.fescar.
  • Компенсационные дела: обычно также называютсяTCC, так как каждая транзакционная операция должна обеспечивать три попытки операции (Try),подтверждать(Confirm) и Компенсация/Отзыв (Cancel), сила непротиворечивости данных ниже, чем у многоэтапной схемы представления, но сложность реализации будет снижена.Очевидный недостаток заключается в том, что для каждой бизнес-транзакции необходимо реализовать три набора операций, а кодов может быть слишком много для компенсационных схем; Кроме того, существует множество сценариев вливания, где TCC не подходит.
  • Сообщение о делах: говорить только здесьRocketMQПроцесс реализации транзакции включает в себя: отправку предварительного сообщения, выполнение локальной транзакции и подтверждение успешной отправки сообщения. Его промежуточное программное обеспечение сообщений хранит сообщения, которые не могут быть успешно обработаны нисходящим потоком, и постоянно пытается передать сообщения потребления нисходящему потоку, в то время как производитель (восходящий) должен предоставитьcheckИнтерфейс для проверки статуса транзакций, которые успешно отправили предварительные сообщения, но не подтвердили окончательный статус отправки сообщения.

Окончательное решение, используемое в проектной практике

Стек технологий личной компании не использует RocketMQ, а в основном использует RabbitMQ, поэтому необходимо адаптировать транзакцию сообщений для RabbitMQ. В настоящее время существует три сценария асинхронного взаимодействия сообщений в бизнес-системе:

  • 1. Отправка сообщений имеет высокую производительность в режиме реального времени и может терпеть потери.
  • 2. Характер отправки сообщений в реальном времени невелик и не может быть потерян.
  • 3. Отправка сообщения имеет высокую производительность в режиме реального времени и не может быть потеряна.

наконец использовали"локальная таблица сообщений"Решение очень простое:

j-t-s-i-a-4.png
j-t-s-i-a-4.png

Основная идея:

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

Псевдокод выглядит следующим образом:

[消息推送实时性高,可以接受丢失-这种情况下可以不需要写入本地消息表 - start]
处理方法(){
    [本地事务开始]
    1、处理业务操作
    [本地事务提交]
    2、组装推送消息并且进行推送
}
[消息推送实时性高,可以接受丢失-这种情况下可以不需要写入本地消息表 - end]

[消息推送实时性低,不能丢失 - start]
处理方法(){
    [本地事务开始]
    1、处理业务操作
    2、组装推送消息并且写入到本地消息表
    [本地事务提交]
}

消息推送调度模块(){
    3、查询本地消息表待推送数据进行推送
}
[消息推送实时性低,不能丢失 - end]

[消息推送实时性高,不能丢失 - start]
处理方法(){
    [本地事务开始]
    1、处理业务操作
    2、组装推送消息并且写入到本地消息表
    [本地事务提交]
    3、消息推送
}

消息推送调度模块(){
    4、查询本地消息表待推送数据进行推送
}
[消息推送实时性高,不能丢失 - end]
  • "Для ситуации «нажатие сообщения в реальном времени велико, потеря приемлема»"На самом деле нет необходимости полагаться на локальную таблицу сообщений, если сообщение собирается и отправляется после отправки транзакции бизнес-операции.В этом случае возникнет проблема потери сообщения из-за недоступности ПО промежуточного слоя очереди сообщений или время простоя локального приложения ("Суть в том, что данные находятся в памяти, непостоянны"), надежность невысокая, но в большинстве случаев проблем нет. При использованииspring-txдекларативные транзакции@Transactionalили программные транзакцииTransactionTemplate,Может"Используйте синхронизатор транзакций для реализации операций RPC, встроенных в блоки кода транзакций бизнес-операций, которые откладываются до тех пор, пока транзакция не будет зафиксирована.", чтобы физическое расположение кода для вызова sub-RPC можно было разместить внутри блока кода транзакции, например:
@Transactional(rollbackFor = RuntimeException.class)
public void process(){
 1.处理业务逻辑
 TransactionSynchronizationManager.getSynchronizations().add(new TransactionSynchronizationAdapter() {
  @Override
  public void afterCommit() {
   2.进行消息推送
  }
 });
}

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

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

Например, структура локальной таблицы сообщений выглядит следующим образом:

CREATE TABLE `t_local_message`(
  id BIGINT PRIMARY KEY COMMENT '主键',
  module INT NOT NULL COMMENT '消息模块',
  tag VARCHAR(20) NOT NULL COMMENT '消息标签',
  business_key VARCHAR(60) NOT NULL COMMENT '业务键',
  queue VARCHAR(60) NOT NULL COMMENT '队列',
  exchange VARCHAR(60) NOT NULL COMMENT '交换器',
  exchange_type VARCHAR(10) NOT NULL COMMENT '交换器类型',
  routing_key VARCHAR(60) NOT NULL COMMENT '路由键',
  retry_times TINYINT NOT NULL DEFAULT 0 COMMENT '重试次数',
  create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建日期时间',
  edit_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '修改日期时间',
  seq_no VARCHAR(60) NOT NULL COMMENT '流水号',
  message_status TINYINT NOT NULL DEFAULT 0 COMMENT '消息状态',
  INDEX idx_business_key(business_key),
  INDEX idx_create_time(create_time),
  UNIQUE uniq_seq_no(seq_no)
)COMMENT '本地消息表';


CREATE TABLE `t_local_message_content`(
  id BIGINT PRIMARY KEY COMMENT '主键',
  message_id BIGINT NOT NULL COMMENT '本地消息表主键',
  message_content TEXT COMMENT '消息内容',
  UNIQUE uniq_message_id(message_id)
)COMMENT '本地消息内容表';

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

Лично лучшая практика для решения распределенных транзакций:

  • "Чтобы избежать реализации распределенных транзакций со строгой согласованностью, основная концепция состоит в том, чтобы отказаться от ACID и перейти к BASE.".
  • Рекомендуется использовать очереди сообщений для развязки между системами.Чтобы обеспечить успех отправки сообщения, отправитель сообщения может независимо прикрепить таблицу сообщений, чтобы связать сообщение, которое нужно отправить, и бизнес-операцию в одной транзакции, и отправить его. асинхронно или по расписанию.
  • Толкатель сообщений (восходящий поток) должен гарантировать, что сообщение правильно доставлено промежуточному программному обеспечению очереди сообщений, а схема потребления или компенсации разрешена потребителем сообщения (нисходящий поток), что будет объяснено в следующей главе.

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

идемпотентное управление

"идемпотент"Термин (идемпотентность) происходит отHTTP/1.1Определения в соглашении:

Методы также могут иметь свойство «идемпотентности» в том смысле, что (помимо ошибок или проблем с истечением срока действия) побочные эффекты N > 0 идентичных запросов такие же, как и для одного запроса.

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

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

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

В настоящее время для идемпотентной обработки на практике используются следующие три аспекта управления:

  • 1. При реализации вызовов идемпотентного интерфейса запись использует распределенные блокировки, используя основныеRedisson, контролировать степень детализации блокировки и время ожидания и удержания блокировки в разумных пределах (промышленность автора требует, чтобы данные были точными, поэтому почти все интерфейсы ядра спроектированы с пессимистичными блокировками,"Я предпочел бы быть медленным, чем ошибаться", на самом деле, если конфликт относительно низок, вы можете рассмотреть возможность использования оптимистической блокировки для оптимизации производительности).
  • 2. Защита от дублирования в бизнес-логике, такой как интерфейс для создания заказа, первый шаг по номеру заказа, чтобы проверить, существует ли уже соответствующий заказ в таблице библиотеки, если он существует, он вернет успех без обработки.
  • 3. Структура таблицы базы данных создает уникальный индекс для логически уникального бизнес-ключа, что является окончательной гарантией на уровне базы данных.

Возьмем пример псевдокода, основанный на идемпотентном управлении потреблением сообщений:

[处理消息消费]
listen(request){
    1、通过业务键构建分布式锁的KEY
    2、通过Redisson构建分布式锁并且加锁
    3、加锁代码中执行业务逻辑(包括去重判断、事务操作和非事务操作等)
    4、finally代码块中释放分布式锁
}

план по компенсации

Схема компенсации в основном включает компенсацию за синхронные HTTP-вызовы и компенсацию за сбои асинхронного потребления сообщений.

Компенсация синхронных вызовов HTTP

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

  • 1. Результат синхронизации приходит в норму, получается конечное состояние, согласованное с нисходящим потоком, и взаимодействие завершается.Принято считать, что успех является конечным состоянием, и никакой компенсации не требуется.
  • 2. Результат синхронизации возвращается нормально, и согласование с нисходящим потоком получено."Нет"Конечное состояние необходимо периодически компенсировать до конечного состояния, иначе будет достигнут верхний предел повторных попыток, который будет помечен как конечное состояние.
  • 3. Результат синхронизации возвращает исключение, наиболее распространенное из которых заключается в том, что нисходящая служба недоступна, а код состояния HTTP — 5XX.

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

Если вы столкнулись со сценарием, связанным с внутренними системами с низким параллелизмомHTTPвзаимодействие, рассмотрите возможность использования"Экспоненциальный откат"Алгоритм повторной попытки, например:

1、第一次调用失败,马上进行第二次重试
2、第二次重试失败,线程休眠2秒
3、第三次重试失败,线程休眠4秒(2^2)
4、第四次重试失败,线程休眠8秒(2^8)
5、第五次重试失败,抛出异常

Если в приведенном выше примере используетсяHystrixТайм-аут управления составляет 1 секунду для переноса HTTP-команды, которая должна быть выполнена для вызова. Описанный выше процесс повторной попытки занимает не более 20 секунд, а взаимодействие между внутренними системами с низким уровнем параллелизма приемлемо.

Однако, если вы столкнулись со сценарием с высоким уровнем параллелизма и высоким приоритетом взаимодействия с пользователем, делать это явно неразумно. Ради безопасности можно принять относительно традиционное и эффективное решение: мгновенное содержимое HTTP-вызова сохраняется в локальной таблице повторов, и эта операция сохранения привязывается к транзакции бизнес-обработки, а"не удалось позвонить"запись, чтобы повторить попытку. Эта схема аналогична схеме, упомянутой выше, для обеспечения успешной отправки сообщения.Вот пример моделирования:

[下单接口请求下游钱包服务扣钱的过程]
process(){
    [事务代码块-start]
    1、处理业务逻辑,保存订单信息,订单状态为扣钱处理中
    2、组装将要向下游钱包服务发起的HTTP调用信息,保存在本地表中
    [事务代码块-end]
    3、事务外进行HTTP调用(OkHttp客户端或者Apache的Http客户端),调用成功更新订单状态为扣钱成功
}

定时调度(){
    4、定时查询订单状态为扣钱处理中的订单进行HTTP调用,调用成功更新订单状态为扣钱成功
}

Компенсация сбоя при асинхронном потреблении сообщений

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

j-t-s-i-a-5.png
j-t-s-i-a-5.png

Если это компенсируется вышестоящими сервисами, возникают две очевидные проблемы:

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

В некоторых недавних проектных практиках определено, что при использовании асинхронного взаимодействия сообщений"Компенсация унифицирована потребителем сообщения". Самый простой способ — использовать метод, аналогичный локальной таблице сообщений, сохранять сообщения, которые не удается обработать, и повторять их. Простая блок-схема выглядит следующим образом:

j-t-s-i-a-6.png
j-t-s-i-a-6.png

Решение для асинхронных сообщений не по порядку

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

Сценарий 1: Вверх по течению службы изменяет гендерную информацию пользователя асинхронно к службе пользователей через очередь сообщений, при условии, что сообщение упрощена следующим образом:

队列:user-service.modify.sex.qeue
消息:
{
   "userId": 长整型,
   "sex": 字符串,可选值是MAN、WOMAN和UNKNOW
}  

Пользовательская служба использует в общей сложности 10 потребительских потоков для мониторинга.user-service.modify.sex.qeueочередь. Предположим, что восходящие сервисы направлены наuser-service.modify.sex.qeueОчередь помещает следующие два сообщения:

第一条消息:
{
   "userId": 1,
   "sex": "MAN"
}  

第二条消息:
{
   "userId": 1,
   "sex": "WOMAN"
}  

Вышеупомянутая отправка сообщений и нисходящая обработка имеют относительно высокую вероятность следующих ситуаций:

j-t-s-i-a-7.png
j-t-s-i-a-7.png

Пользователь с исходным идентификатором пользователя 1 сначала изменил пол на МУЖЧИНА (первый запрос), а затем изменил его на ЖЕНЩИНА (второй запрос) и, наконец, увидел, что обновленный пол может быть МУЖЧИНА, что явно неразумно. Проблема, которую хочет проиллюстрировать этот необоснованный пример, заключается в следующем: из-за асинхронного взаимодействия сообщений время обработки сообщений нижестоящими службами может не соответствовать времени отправки сообщений вверх по течению, что может привести к нарушению бизнес-статуса. Для решения этой проблемы есть несколько возможных идей:

  • Вариант 1. Если требование параллелизма невелико, вы можете в полной мере использовать очередь сообщений.FIFOхарактеристики (этоRabbitMQРеализовано, другое промежуточное ПО очереди сообщений не определено), установите для потока потребителя нижестоящей службы значение 1, тогда синхронизация push-сообщения восходящего потока и сообщения потребления нижестоящего потока будет согласована.
  • Вариант 2. Используйте HTTP-вызовы, для этого требуется взаимодействие внешнего интерфейса или клиента приложения, а запрос может быть последовательным.

Сценарий 2: асинхронная обработка сообщений без требований к времени, но для окончательного отображения требуется время. Это может быть немного абстрактно, например: при заимствовании 10 000 юаней на Borrow пользователь погашает их кратно (например, план погашения 1: 2000, 3000, 5000; план погашения 2: 1000, 1000, 1000, 7000 и т. д. .), погашение каждый раз разное, и окончательный счет требуется отображать в соответствии с порядком погашения пользователя.

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

Решение простое: прикрепите флаг с возрастающей или убывающей тенденцией при отправке сообщения, например, с помощью флага с отметкой времени или с помощьюSnowflakeАлгоритм генерирует самоувеличивающееся длинное целое число в качестве порядкового номера, а затем сортирует его в соответствии с порядковым номером, чтобы получить последовательность операций над сообщением (порядковый номер необходимо сохранить в дальнейшем), но фактическая обработка сообщения не требуется. воспринимать время сообщения.

Асинхронный обмен сообщениями в сочетании с управлением состоянием

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

j-t-s-i-a-8.png
j-t-s-i-a-8.png

Решение этой проблемы заключается в использовании синхронных вызовов (на самом деле что-то вродеTCC,2PCили3PCи т. д. по существу являются синхронными вызовами), и строгая согласованность может быть достигнута при условии допущения потери производительности. В этом разделе не обсуждается, как это сделать в случае синхронных вызовов, а основное внимание уделяется тому, как использовать очереди сообщений изBASEугол «для достижения относительно высокой степени согласованности». Чтобы сначала абстрагироваться от этого примера, предположим, что таблицы счетов обеих систем спроектированы следующим образом:

CREATE TABLE `t_account`(
    id BIGINT PRIMARY KEY COMMENT '主键',
    user_id BIGINT NOT NULL COMMENT '用户ID',
    balance DECIMAL(10,2) NOT NULL DEFAULT 0 COMMENT '账户余额',
    create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
    edit_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '修改时间',
    version BIGINT NOT NULL DEFAULT 0 COMMENT '版本'
    // 省略索引
)COMMENT '账户表';

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

[A系统本地事务-start]
1、A系统t_account表X用户余额减去1000
2、A系统流水表写入一条用户X的预扣1000的记录,标记状态为处理中,生成全局唯一的流水号记为SEQ_NO
[A系统本地事务-end]
3、A系统通过消息队列推送一条用户X扣减1000的消息(一定要附带流水号SEQ_NO)到消息队列中间件(这里可以用上文提到的技巧确保消息推送成功)
[B系统本地事务-start]
4、B系统t_account表X用户余额加上1000
5、B系统流水表写入一条用户X的余额变更(增加)1000的记录 <= 注意这里B系统的流水只能insert不能update
[B系统本地事务-end]
6、B系统推送处理X用户余额处理成功的消息到消息队列中间件,一定要附带流水号SEQ_NO(这里可以用上文提到的技巧确保消息推送成功)
[A系统本地事务-start]
7、A系统更新流水表中X用户流水号为SEQ_NO的预扣记录的状态为处理成功(这一步一定要做好幂等控制,可以考虑用SEQ_NO作为分布式锁的KEY)
[A系统本地事务-end]

其他:
[A系统流水表处理中的记录需要定时轮询和重试]
1、定时调度重试A系统流水表中状态为处理中的记录

[A-B系统日切对账模块]
1、日切,用A系统中处理成功的T-1日流水记录和B系统中的流水表所有T-1日的记录进行对账
j-t-s-i-a-9.png
j-t-s-i-a-9.png

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

резюме

Вы обнаружите, что в статье используется много схем."Отложенное содержимое записывается в локальную таблицу + триггер в реальном времени вне транзакции + компенсация планирования времени"Этот режим, на самом деле, что я хочу выразить, так это то, что этот режим является относительно распространенным режимом в текущем распределенном решении, которое может в основном соответствовать обработке различных сложных сценариев, таких как распределенные транзакции, синхронная и асинхронная компенсация и в режиме реального времени срабатывание не в реальном времени. Есть также некоторые очевидные проблемы с этим шаблоном (которые обычно встречаются при практике):

  • 1. Неразумный дизайн или неразумная обработка библиотечной таблицы (локальной таблицы сообщений) может легко стать узким местом базы данных.
  • 2. Логический код обработки компенсации или хранения локальной таблицы легко может быть избыточным и поврежденным.
  • 3. В крайних случаях сценарий аварийного восстановления может вывести из строя скрытую опасность сервиса.

На самом деле чаще его нужно анализировать в сочетании с существующими системами или сценариями, а последующую оптимизацию проводить через мониторинг и анализ данных. после всего,"Архитектура итеративна, а не спроектирована". Кроме того, упомянутый текст"Окончательное решение фактически написал еще одну статью для подробного анализа:"Повторно используемая схема сообщений транзакций на основе RabbitMQ.

(Конец этой статьи e-a-20190323 c-14-d 996 Это статья, написанная в конце марта 2019 года, и я надеюсь, что она не устарела сейчас)