【问题标题】:How to save message into database and send response into topic eventually consistent?如何将消息保存到数据库并将响应发送到主题最终一致?
【发布时间】:2019-11-14 00:46:21
【问题描述】:

我有以下 rabbitMq 消费者:

Consumer consumer = new DefaultConsumer(channel) {
    @Override
     public void handleDelivery(String consumerTag, Envelope envelope, MQP.BasicProperties properties, byte[] body) throws IOException {
            String message = new String(body, "UTF-8");
            sendNotificationIntoTopic(message);
            saveIntoDatabase(message);
     }
};

可能会出现以下情况:

  1. 消息已成功发送到主题
  2. 与数据库的连接丢失,因此数据库插入失败。

因此我们的数据不一致。

预期结果要么两个动作都成功执行,要么都没有执行。

任何解决方案我该如何实现它?

附言

目前我有以下想法(请评论)

我们可以假设代理不会丢失任何消息。

我们必须订阅我们想要发送的主题。

  1. 将条目保存到数据库并设置字段status 值为“待定”
  2. 尝试向主题发送数据。如果发送成功 - 更新字段 status 值为“成功”
  3. 我们必须有一个已调度的作业,它必须检查具有待处理状态的行。目前可能有两种情况:
    3.1 根本没有发送通知
    3.2 通知已发送但保存到数据库失败(概率很低但有可能)

    所以我们必须以某种方式区分这两种情况:我们可以将来自主题的消息存储在集合中,并且作业可以检查消息是否被接受。因此,如果作业找到与数据库行对应的消息,我们必须将状态更新为“成功”。否则我们必须从数据库中删除条目。

我认为我的想法有一些弱点(例如,如果我们有多节点应用程序,我们必须将消息存储在 hazelcast(或类似物)中,但这是假设失败的额外点)

【问题讨论】:

  • @user7294900 我们的重试次数有限。如果代理关闭,我们可以用尽所有尝试,并且我们再次遇到数据不一致
  • @user7294900 我不知道)但是我没有遇到过 10 亿次重试尝试的系统
  • 解决方案是使用支持 JMS 和 XA 事务的消息系统,并使用 XA 事务管理器。或者拥有能够容忍不一致的业务逻辑。
  • @JB Nizet 当有人听到有关 XA 交易的消息时,他通常会变得紧张)
  • @JB Nizet 听起来很有趣,但我无法想象该怎么做

标签: java transactions rabbitmq messagebroker eventual-consistency


【解决方案1】:

这是一个尝试取消确认模式https://servicecomb.apache.org/docs/distributed_saga_3/的示例,它应该能够处理您的问题。您应该容忍通过队列重复提交数据的一些机会。这是一个例子:

  1. 定义抽象操作并为操作分配 ID 和时间戳。
  2. 将状态 Pending 写入数据库(您可以在与 1 相同的步骤中执行此操作)
  3. 编写一个侦听器,轮询数据库中状态为挂起且早于“超时”的所有操作
  4. 对于每个挂起的操作,通过具有分配 ID 的队列发送数据。
  5. 接收方应该知道 ID,如果 ID 已被处理,则不会发生任何事情。

6A。如果您需要 100% 完成操作,您需要第二个队列,接收方将在其中发布消息 ID - DONE。如果不需要这种一致性,请跳过此步骤。或者,它可以发布 ID -Failed 失败原因。

6B。提交方要么通过将状态 DONE 写入数据库来等待来自 6A 的消息完成操作。

  • 一旦 sertine 超时或某个重试限制已经过去。您将状态写入操作 FAIL。
  • 您可以通过 ID 回滚向接收方操作发送消息。

请注意,所有这些步骤都不涉及技术交易。您可以使用非事务性数据库来执行此操作。

我所写的是尝试取消确认模式的变体,其中每个消息接收者都应该知道如何管理自己的数据。

【讨论】:

  • @gstackoverflow 它更详细。您没有考虑到您可能会发布到队列,但由于应用程序错误或验证规则或其他原因,消息仍然没有传递给收件人,或者正在传递但未处理... que 并不一定意味着成功。
  • @gstackoverflow 您也没有考虑到发生错误时可能发生的各种补偿方案。
  • 工作是自我补偿的。出现故障会重启
  • @gstackoverflow 不是。在队列上成功发布并不意味着没有错误。您最终可能会遇到消息已被消费但接收方的操作尚未以良好方式完成的情况。你将如何捕捉这个?按照您的 3 个步骤,您将在成功发布消息后写入数据库 Success。您的算法不能保证收件人和发件人的状态一致。
  • 可能有话题。我可能不想检查所有主题订阅者是否能够处理消息
【解决方案2】:
  1. 在侦听器中保存数据库行,字段 staus='pending'
  2. 另一个作业(分离的线程)将从数据库中获取所有待处理的行,并为每行获取以下信息:
    2.1 向topic发送数据
    2.2 存入数据库

如果我们在第 1 步失败 - 一切正常 - 数据处于一致状态,因为作业不会知道有关该数据的任何信息

如果我们在步骤 2.1 上失败了 - 没问题,下一个作业调用将尝试处理它

如果我们在步骤 2.2 上失败了 - 如果我们在这里失败了 - 这意味着下一次作业调用将再次处理相同的数据。乍一看,您可能会认为这是一个问题。但是你的消费者必须是幂等的——这意味着它必须理解消息已经被处理并跳过处理。此要求是所有消息代理都保证消息将至少传递一次的结果。因此,无论如何,我们的消费者都必须为重复消息做好准备。又没问题了。

【讨论】:

    【解决方案3】:

    这是我如何做的伪代码:(假设 dao 层具有事务能力,而您的消息传递层没有)

        //Start a transaction
        try {
                    String message = new String(body, "UTF-8");
                   // Ordering is important here as I'm assuming the database has commit and rollback capabilities, but the messaging system doesnt. 
                    saveIntoDatabase(message);
                    sendNotificationIntoTopic(message);
    
        } catch (MessageDeliveryException) {
            // rollback the transaction
            // Throw a domain specific exception
        }
       //commit the transaction
    

    场景:
    1.如果数据库失败,消息将不会发送,因为异常会破坏代码流。
    2.如果数据库调用成功,消息系统下发失败,捕获异常并回滚数据库更改

    记录和重放失败所需的所有操作都可以在此方法之外

    【讨论】:

    • 数据库在回滚期间可能失败
    • @gstackoverflow 在这种情况下,事务还没有提交,所以它是安全的
    【解决方案4】:

    如果有足够的时间修改设计,建议使用类似 JTA 的 API 来管理 2phase 提交。甚至 weblogic 和 WebSphere 也支持 XA 资源进行 2 阶段提交。

    如果时间线比较短,建议按照下面的方法来减少失败的差距。

    • 发送数据主题(无提交)(如果主题关闭,请间隔重试)
    • 将数据写入数据库
    • 提交数据库
    • 提交主题

    这里只有在第 4 步失败时才会发生失败。这将导致再次发送相同的消息。所以接收系统会收到重复的消息。在 JMS2.0 结构中,每条消息都有唯一的 messageID 和 CorrelationID。所以查找重复有点简单(但这要在接收系统处理)

    这两种情况也适用于集群环境。


    严格按照您的情况,认为以下步骤可能有助于解决您的问题

    为您的主题订阅一个 listener listener-1。

    进程-1

    • 为消息 msg-1 添加状态为“待发送”的 DB 条目
    • 向主题发送消息 msg-1。在任何主题失败的情况下重试发送 如果在某些重试后第 2 步失败,则 process-1 必须在发送任何新消息之前重新发送 msg-1 或第 1 步回滚

    监听器 1

    • 使用订阅的侦听器,从 Topic 中读取引用(meesageID/correlationID),并将 DB 状态更新为 SENT,并从 topic 中读取/删除消息。万一引用读取成功并且数据库更新失败,主题仍然有消息。所以下一次读取将更新数据库。 Incase 数据库更新成功并且消息删除失败。侦听器将再次读取并尝试更新已经完成的消息。所以验证后可以忽略。

    如果侦听器本身关闭,主题将有消息,直到侦听器读取消息。在此之前,SENT 消息将处于“待发送”状态。

    【讨论】:

    • 如果我们不能像您提到的那样提交主题,它将无法工作。不能算是解决办法
    猜你喜欢
    • 2020-05-09
    • 2019-10-17
    • 2012-08-30
    • 2021-03-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-08-30
    相关资源
    最近更新 更多