【问题标题】:Google NDB Datastore Checking Account / Wallet for Users. How to calculate balance用户的 Google NDB 数据存储检查帐户/电子钱包。如何计算余额
【发布时间】:2018-02-12 00:49:52
【问题描述】:

基本上,使用 Google Cloud Datastore 计算滚动余额以确定何时补充用户钱包的最佳方法是什么?

付款和交易

我维护一个支付平台,我们的用户可以在该平台上从第三方机构购买各种物品。不幸的是,我所在行业的性质是,这些支付事件不是实时发送给我们的,它们会打包成一批并在几小时到几周后发送给我们。

这些是影响用户钱包余额的主要对象:

class Transaction(ndb.model):
    user = ndb.KeyProperty(User, required=True)
    amount = ndb.FloatProperty(required=True)
    # ... other fields

class Payment(ndb.model):
    user = ndb.KeyProperty(User, required=True)
    amount = ndb.FloatProperty(required=True)
    # ... other fields

    @classmethod
    def charge(cls, user, amount):
        # ... make a call to braintree/stripe & save result if successful

(未显示退款、“商店信用”、调整等)

钱包

但是,大部分交易金额小于 1 美元。由于我们必须将信用卡处理成本转嫁给用户,因此我们的用户与我们一起维护钱包以尽量减少这些费用。

他们可以充值 10 美元到 200 美元,交易会从余额中扣除,当他们的余额较低(低于 2 美元)时,我们会从他们的卡中扣款以补充他们的帐户。

这就是我设想的钱包活动模型的工作方式

class WalletActivity(ndb.Model):
    user = ndb.KeyProperty(User, required=True)
    post_date = ndb.DateTimeProperty(required=True)
    balance_increment = ndb.FloatProperty(required=True)
    balance_result = ndb.FloatProperty(required=True)
    # the key to the Transaction or Payment object that this is for
    object_key = ndb.KeyProperty(required=True)

    @classmethod
    def create(cls, obj, previous_balance):
        return WalletActivity(
            user_key=obj.user,
            post_date=datetime.datetime.now(),
            balance_increment=obj.amount,
            balance_result=previous_balance+obj.amount,
            object_key=obj.key)

    @classmethod
    def fetch_last_wallet_activity(cls, user_key):
        return cls.query(cls.user == user_key).order(-cls.post_date).get()

计算余额

为了确定平衡,频谱的两端似乎是:

  • 即时计算,汇总帐户的整个钱包历史记录
  • 存储预计算值 (WalletActivity.fetch_last_wallet_activity().balance_result)

这里的正确答案听起来像是 2 的组合。 在每个帐户的每一天结束时存储某种 BalanceUpdate / WalletDaySummary 对象。 然后,您只需总结今天的活动并将其添加到昨天的 BalanceUpdate。 https://stackoverflow.com/a/4376221/4458510

class BalanceUpdate(ndb.model):
    user = ndb.KeyProperty(User)
    cut_off_date = ndb.DateTimeProperty()
    balance = ndb.IntegerProperty()

    @classmethod
    def current_balance(cls, user_key):
        last_balance_update = cls.query(cls.user == user_key).order(
            -cls.cut_off_date).get()
        recent_wallet_activity = WalletActivity.query(cls.user == user_key, 
            cls.post_date > last_balance_update.cut_off_date).fetch()
        return (last_balance_update.balance + 
            sum([i.balance_increment for i in recent_wallet_activity]))

但是,这可能不适用于在一天内产生大量交易的公司帐户。 使用最近的WalletActivity的balance_result可能会更好

如何处理交易

选项 1

要处理一批交易,我们会

  1. 获取用户余额
  2. 如果现有余额不足,请为他们的帐户充值
  3. 将交易添加到他们的钱包中

代码:

def _process_transactions(user, transactions, last_wallet_activity):
    transactions_amount = sum([i.amount for i in transactions])
    # 2. Replenish their account if the existing balance is low
    if last_wallet_activity.balance_result - transactions_amount < user.wallet_bottom_threshold:
        payment = Payment.charge(
            user=user,
            amount=user.wallet_replenish_amount + transactions_amount)
        payment.put()
        last_wallet_activity = WalletActivity.create(
            obj=payment,
            previous_balance=last_wallet_activity.balance_result)
        last_wallet_activity.put()
    # 3. Add the transactions to their wallet
    new_objects = []
    for transaction in transactions:
        last_wallet_activity = WalletActivity.create(
            obj=transaction,
            previous_balance=last_wallet_activity.balance_result)
        new_objects.append(last_wallet_activity)
    ndb.put_multi(new_objects)
    return new_objects

def process_transactions_1(user, transactions):
    # 1. Get the user's balance from the last WalletActivity
    last_wallet_activity = WalletActivity.fetch_last_wallet_activity(user_key=user.key)
    return _process_transactions(user, transactions, last_wallet_activity)

WalletActivity.fetch_last_wallet_activity().balance_result 和 BalanceUpdate.current_balance() 是数据存储区查询最终是一致的。

我考虑过使用实体组和祖先查询,但听起来您会遇到争用错误:

选项 2 - 按键获取 Last WalletActivity

我们可以追踪最后一个WalletActivity的key,因为key fetching是强一致性的:

class LastWalletActivity(ndb.Model):
    last_wallet_activity = ndb.KeyProperty(WalletActivity, required=True)

    @classmethod
    def get_for_user(cls, user_key):
        # LastWalletActivity has the same key as the user it is for
        return ndb.Key(cls, user_key.id()).get(use_cache=False, use_memcache=False)

def process_transactions_2(user, transactions):
    # 1. Get the user's balance from the last WalletActivity
    last_wallet_activity = LastWalletActivity.get_for_user(user_key=user.key)
    new_objects = _process_transactions(user, transactions, last_wallet_activity.last_wallet_activity)

    # update LastWalletActivity
    last_wallet_activity.last_wallet_activity = new_objects[-1].key
    last_wallet_activity.put()
    return new_objects

或者,我可以将last_wallet_activity 存储在User 对象上,但我不想担心竞争条件 用户更新他们的电子邮件并清除我对last_wallet_activity 的新值

选项 3 - 付款锁定

但是如果竞争条件有 2 个作业试图同时处理同一个用户的事务呢? 我们可以添加另一个对象来“锁定”一个帐户。

class UserPaymentLock(ndb.Model):
    lock_time = ndb.DateTimeProperty(auto_now_add=True)

    @classmethod
    @ndb.transactional()
    def lock_user(cls, user_key):
        # UserPaymentLock has the same key as the user it is for
        key = ndb.Key(cls, user_key.id())
        lock = key.get(use_cache=False, use_memcache=False)
        if lock:
            # If the lock is older than a minute, still return False, but delete it
            # There are situations where the instance can crash and a user may never get unlocked
            if datetime.datetime.now() - lock.lock_time > datetime.timedelta(seconds=60):
                lock.key.delete()
            return False
        key.put()
        return True

    @classmethod
    def unlock_user(cls, user_key):
        ndb.Key(cls, user_key.id()).delete()

def process_transactions_3(user, transactions):
    # Attempt to lock the account, abort & try again if already locked 
    if not UserPaymentLock.lock_user(user_key=user.key):
        raise Exception("Unable to acquire payment lock")

    # 1. Get the user's balance from the last WalletActivity
    last_wallet_activity = LastWalletActivity.get_for_user(user_key=user.key)
    new_objects = _process_transactions(user, transactions, last_wallet_activity.last_wallet_activity)

    # update LastWalletActivity
    last_wallet_activity.last_wallet_activity = new_objects[-1].key
    last_wallet_activity.put()

    # unlock the account
    UserPaymentLock.unlock_user(user_key=user.key)
    return new_objects

我曾想过尝试将整个事情包在一个事务中,但我需要防止对 Braintree/stripe 进行 2 个 http。

我倾向于选项 3,但随着我引入的每个新模型,系统感觉越来越脆弱。

【问题讨论】:

  • 如果每秒有超过 1 个写入操作进入同一实体组,则会发生争用错误。这也适用于从 Datastore 事务中的 Datastore 读取而没有显式写回的实体(我的理解是它们也是在后台写入以进行序列化的)。由于我没有看到特定的面向用户的请求或时间限制,如果您将所有事务和用户加载到同一个实体组 + 使用 Datastore 事务,您可以限制写入操作并实现安全失败和每个用户的指数回退。
  • PS:一般来说,设计应用程序以在写入限制内安全工作的方式比编写尝试实现事务方面的逻辑更容易。
  • 我记得你最初的问题,我同意这个答案,特别是因为你提到了拥有数百万笔交易的公司账户,假设它们是实时处理的。但是,在您上面的问题中,显然数小时甚至数周都无关紧要,因此您可以将它们分批排队和批处理。每个账户每秒 1 个 Datastore 事务,最多 500 个金融事务,相当于每个账户每小时 180 万。即使有重试和中断,对于企业帐户来说,这可能就足够了吗?
  • @RodrigoC。我已经发布了我的信息作为答案,还添加了一些其他注意事项和建议。希望这很有用。
  • @Ani,确实如此。再次感谢您为此以及您在这个社区所做的所有辛勤工作。

标签: python google-app-engine google-cloud-datastore


【解决方案1】:

关于争用的设计注意事项

总的来说,我完全同意 Dan's answer 和您的 original question,尽管在您的特定用例中,使用大型实体组可能是合理的。

对于同一实体组,每秒写入操作超过 1 次时可能会发生争用错误,例如一个特定的钱包。此限制也适用于在 Datastore 事务中从 Datastore 读取而没有显式写回的实体(我的理解是它们是 written together with the modified entities for serializability)。

虽然 1 秒规则不是强制限制,但根据我的经验,Cloud Datastore 通常可以处理略高于此限制的短突发,但无法保证,通常建议和最佳实践避免大型实体组,特别是在写操作不是来自同一用户的情况下。相比之下,将用户发布的所有 cmets 存储在同一实体组(作者 = 父级)中可能是安全的,因为用户每秒发布超过一条评论的可能性很小,甚至可能是不可取的。另一个示例可能是对时间不敏感且不面向用户的后台任务,其中实体组的写入操作是按实体组编排的,或者如果发生争用,至少可以显着回退写入操作。

在突然以非常高的速率添加新实体的情况下,争用错误也可能由单调增加的键/ID 或索引属性和复合索引引起,其中索引值彼此太接近(例如时间戳)。建议让 Datastore 自动创建新实体的 ID(例如用户 ID),因为 Datastore 会将密钥分散得足够远。并且要么避免对可能出现单调递增值的属性进行索引,要么在值前面加上一个哈希值。

Cloud Datastore 文章 Best Practices 包含一个关于 Designing for scale 的部分,该部分提供了非常有用的建议。

也就是说,与编写尝试模仿的应用程序逻辑相比,设计应用程序以在写入限制内安全工作并依赖 Datastore(或其他数据库)的事务性和强一致性支持的方式更容易交易方面。竞争条件、死锁以及更多可能出错的因素,使系统更加脆弱并容易出错。

(A) 小实体组中的钱包活动

在最初的问题中,您提到了公司帐户和付款,这提出了一些实时付款解决方案。场景:该公司的数千名用户可以为单个帐户提交新交易,但许多人可能同时这样做的风险非常高。如果每笔交易都针对同一个实体组(公司账户),这很容易导致争用错误。如果您在写入操作中实现重试,则会导致延迟延长,直到用户得到对其事务请求的响应。但即使重试,写操作也很可能经常失败,用户会经常面临服务器错误。这会带来糟糕的用户体验。

我倾向于选项 (2),但使用 Wallet 类型。您表示担心将last_wallet_activity 存储在哪里。您可以拥有自己的Wallet 类型,它始终具有与User 相同的ID。在这种情况下,您可以拥有两个单独的实体组,并且不关心用户触发的 User 对象的中间更改。我也会使用 Datastore 事务。它将允许在同一事务中最多有 25 个不同的实体组,即一批中可能有 23 个事件。使用您当前的设计,这也将是您每个钱包的最大写入率。不确定这是您的应用可接受的限制。

但是,如果某个钱包确实有很多交易(通过频繁更新last_wallet_activity),您可能会再次面临争用的风险。为了获得更高的每个钱包写入率并避免此类争用,您还可以将选项 (1) 与某种 sharding 结合使用。

选项 (3) 尝试实现事务方面(见上文)。选项 (1) 确实存在查询仅是最终一致的事实。

(B) 单个大型实体组中的钱包活动

但是,这里的问题提到所有这些钱包活动都是在后台处理的(不是由用户请求),并且这些事件不是实时处理的,而是延迟数小时或数周。这将允许通过后台任务批量处理它们。假设您的应用程序将保留在另一个 Cloud Datastore limits 内,例如对于 10 MiB 的事务的最大大小,您的应用程序可以为每个公司钱包(实体组)每秒将多达 500 个钱包活动批处理到单个写入操作中。这相当于每个公司账户每小时 180 万次钱包活动。即使有重试和中断,这对于最大的公司帐户来说是否足够?如果是,并且如果您的产品永远不会更改为面向用户的实时钱包活动,我不明白为什么您不应该将钱包活动放入每个钱包的实体组中。当然,另一种方法是每个公司帐户只拥有多个钱包。

在这种情况下,您的选项 (1) 应该可以工作,因为祖先查询(其中钱包是 WalletActivity 查询中的祖先)是高度一致的。

使用ndb.put_multi()

我看到你在同一个请求中使用了多个put() 调用,你可以完美地收集一些toPut 列表中的实体,然后将它们全部写在一起。通过这样做,您可以节省实例运行时间并减少对同一实体组的写入操作数。

如何避免重复支付请求

关于向 Braintree 或 Stripe 的付款请求:

  1. 在将钱包活动添加到将写入 Datastore 的下一批之前,请检查钱包的余额是否足够。
  2. 如果余额不足,请停止添加更多活动。
  3. 在您将批处理写入数据存储的数据存储事务中,添加transactional push task(如果事务失败,则不会创建)。我相信,GAE/NDB 每个 HTTP 请求最多接受 5 个事务任务。
  4. 此任务将负责将请求发送到 Braintree/Stripe,并更新钱包的余额。这也应该在 Datastore 事务中,具有事务任务以继续处理之前离开的事件。

您需要处理 Braintree/Stripe 拒绝付款请求的情况。

我也不知道你的事件是如何到达应用程序的,所以我不确定协调每个钱包的批处理任务的最佳方式。上面的模式表明,对于每个钱包,只有一个任务在运行(即,同一个钱包不是并行的多个任务/批次)。但是有不同的方法可以做到这一点,具体取决于事件到达您的应用的方式。

【讨论】:

  • 只是一个注释 - 争用不仅限于写入同一个实体组。我一开始也是这么想的,结果被烧死了,不得不做一些重大的重新设计来避免它。见stackoverflow.com/questions/32927068/…。
  • 感谢@DanCornilescu(顺便说一句,也感谢您对 SO 的所有贡献)。只是为了澄清您的注释与我在第二段中链接的可序列化部分不同:也适用于从不写回该组或事务中任何其他组的事务?对于可序列化部分,我知道一旦任何实体组发生写入,事务的所有实体组都会得到写入,即使事务未修改也是如此。您的回答表明“标记”导致了争用。我的理解正确吗?
  • 我不知道确切的技术细节,我只是推断它一定是标记旁边的东西,但是是的,这就是我的意思 - 事务在读取操作时失败,它们看起来像这样:@987654341 @,键指向一个(配置)实体,该实体不是由当时执行的任何操作写入的。
  • 是的,对于作为服务配置工作的实体,这与我遇到的相同错误,并且仅在事务内部读取,而写入发生在其他实体组中。我记得几年前一位 Google 员工向我解释说 Cloud Datastore 使用乐观并发,但我不记得在哪里获得了这些信息。但是,我在这里找到了这篇较早的文章:cloud.google.com/appengine/articles/scaling/…,这让我相信完全只读的事务不会发生争用。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2014-05-05
  • 2020-10-02
  • 2017-10-17
  • 2011-05-21
  • 2011-09-14
  • 2021-08-04
  • 2022-10-17
相关资源
最近更新 更多