【发布时间】: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
要处理一批交易,我们会
- 获取用户余额
- 如果现有余额不足,请为他们的帐户充值
- 将交易添加到他们的钱包中
代码:
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