【发布时间】:2017-07-01 10:42:11
【问题描述】:
我有以下典型场景:
- 用于购买产品的订单服务。充当分布式事务的指挥官。
- 包含产品列表及其库存的产品服务。
-
一种支付服务。
Orders DB Products DB | | --------------- ---------------- ---------------- | OrderService | | ProductService | | PaymentService | --------------- ---------------- ---------------- | | | | -------------------- | --------------- | Kafka orders topic |------------- ---------------------
正常流程是:
- 用户订购产品。
- 订单服务在 DB 中创建订单并在 Kafka 主题“订单”中发布消息以预留产品 (PRODUCT_RESERVE_REQUEST)。
- 产品服务将其数据库中的产品库存减少一个单位,并在“订单”中发布一条消息说 PRODUCT_RESERVED
- Order 服务获取 PRODUCT_RESERVED 消息并命令支付发布消息 PAYMENT_REQUESTED
- 支付服务订购付款并回复一条消息已支付
- 订单服务读取 PAYED 消息并将订单标记为 COMPLETED,完成交易。
我在处理错误情况时遇到了麻烦,例如:让我们假设:
- 支付服务未能为产品收费,因此它发布消息PAYMENT_FAILED
- 订单服务响应发布消息 UNDO_PRODUCT_RESERVATION
- 产品服务增加数据库中的库存以取消预订并发布 PRODUCT_UNRESERVATION_COMPLETED
- 订单服务完成交易,将订单的最终状态保存为 CANCELLED_PAYMENT_FAILED。
在这种情况下,假设无论出于何种原因,订单服务发布了 UNDO_PRODUCT_RESERVATION 消息但没有收到 PRODUCT_UNRESERVATION_COMPLETED 消息,因此它重新尝试发布另一个 UNDO_PRODUCT_RESERVATION 消息。
现在,假设同一订单的这两条 UNDO_PRODUCT_RESERVATION 消息最终到达 ProductService。如果我同时处理它们,我最终可能会为产品设置无效库存。
在这种情况下如何实现幂等性?
更新:
按照 Artem 的说明,我现在可以检测到重复的消息(通过检查消息标头)并忽略它们,但可能仍然存在以下情况,我不应该忽略重复的消息:
- Order Service 发送 UNDO_PRODUCT_RESERVATION
- 产品服务收到消息并开始处理它,但在更新库存之前崩溃。
- Order Service 未收到响应,因此它重新尝试发送 UNDO_PRODUCT_RESERVATION
- 产品服务知道这是一条重复的消息但是,在这种情况下它应该再次重复处理。
你能帮我想出一种方法来支持这种情况吗?我如何区分何时应该丢弃消息或重新处理它?
【问题讨论】:
-
嗨,Artem,它看起来不错...我正在使用 Spring Cloud Stream,因此对配置提出了一些疑问。是否像添加spring集成依赖和配置
IdempotentReceiverInterceptor一样简单?我正在使用@StreamListener(OrderProcessor.INPUT)来注释处理方法。@IdempotentReceiver("idempotentReceiverInterceptor")注释可以在他们身上正常工作吗? -
嗯,SCSt 完全基于 Spring Integration。我不确定谁说你不一样。那里已经有很多适当的依赖关系。不,这不适用于该注释。但是你可以切换到
@ServiceActivator,消息的结果相同,但幂等接收器将在这里工作。 -
太棒了!我要试一试。谢谢!
-
我修改了代码,可以看到应用了拦截器。但是,默认情况下,消息具有不同的 id (UUID)。我尝试手动设置一个,但它不起作用:
MessageBuilder.withPayload(msg.getMessage().getBytes()).setHeader(MessageHeaders.ID, msg.getId()+"-"+msg.getState()).build();异常:java.lang.IllegalArgumentException: 'id' header is read-only。如何提供生成的唯一 ID?问候
标签: spring-cloud microservices distributed-transactions spring-cloud-stream spring-kafka