多条件下大额DeFi流动性头寸手续费分配高效实现问询
百万级DataFrame下流动性头寸手续费分配的高效实现
问题背景
我有两个百万级规模的DataFrame:
test_swaps:记录各交易池内的所有交易及对应手续费,结构示例:
BLOCK_TIMESTAMP LIQUIDITY POOL_ADDRESS PRICE_0_1 FEES 0 2021-12-02 03:42:56+00:00 7.770303e+21 0x360b9726186c0f62cc719450685ce70280774dc8 215.174737 163.787000 1 2021-12-02 04:02:18+00:00 7.770303e+21 0x360b9726186c0f62cc719450685ce70280774dc8 209.796784 223.138500 2 2021-12-02 04:03:52+00:00 7.770303e+21 0x360b9726186c0f62cc719450685ce70280774dc8 206.188879 199.961300 3 2021-12-02 06:55:37+00:00 8.165560e+19 0xfaa318479b7755b2dbfdd34dc306cb28b420ad12 203.100999 0.044125 4 2021-12-02 04:09:03+00:00 7.770303e+21 0x360b9726186c0f62cc719450685ce70280774dc8 204.329947 22.049300 ...
test_actions:按NFT代币ID记录各交易池的流动性头寸变动,结构示例:
BLOCK_TIMESTAMP NF_TOKEN_ID LIQUIDITY_cumsum POOL_ADDRESS PRICE_LOWER_0_1 PRICE_UPPER_0_1 0 2021-05-05 19:31:28+00:00 374.0 2.629662e+20 0xfaa318479b7755b2dbfdd34dc306cb28b420ad12 79.820552 81.189032 1 2021-05-05 21:12:56+00:00 374.0 0.000000e+00 0xfaa318479b7755b2dbfdd34dc306cb28b420ad12 79.820552 81.189032 2 2021-05-05 20:03:50+00:00 539.0 7.412937e+23 0x360b9726186c0f62cc719450685ce70280774dc8 0.012037 0.012781 3 2021-05-05 20:13:27+00:00 539.0 0.000000e+00 0x360b9726186c0f62cc719450685ce70280774dc8 0.012037 0.012781 4 2021-05-05 20:29:05+00:00 636.0 4.235670e+19 0x360b9726186c0f62cc719450685ce70280774dc8 66.672329 95.561691 ...
目标与规则
生成包含NF_TOKEN_ID、POOL_ADDRESS、FEES_ACCRUED的结果表,展示每个流动性头寸累计的手续费。头寸累计手续费需满足以下条件:
- 头寸与交易属于同一
POOL_ADDRESS; - 交易价格
PRICE_0_1落在头寸的价格区间[PRICE_LOWER_0_1, PRICE_UPPER_0_1]内; - 交易发生时头寸处于活跃状态:
- 每个
NF_TOKEN_ID的记录按时间排序,首次记录为激活时间; - 若最后一条记录的
LIQUIDITY_cumsum为0,头寸活跃时间为「激活时间到最后一条记录时间」; - 若最后一条记录的
LIQUIDITY_cumsum非0,头寸持续活跃至所有交易结束。
- 每个
手续费按头寸流动性占比分配:头寸LIQUIDITY_cumsum / 交易发生时池内总LIQUIDITY,再乘以该笔交易的FEES。
原代码的问题
原代码采用双重循环遍历所有交易和头寸,且每次循环内多次查询DataFrame,时间复杂度为O(M*N)(M为交易数,N为头寸数),对于百万级数据完全不可行,导致运行极慢甚至报错。
优化方案
步骤1:预处理头寸数据,生成活跃时间区间
先整理每个头寸的有效活跃时间段、流动性数值和价格区间,避免重复计算:
import pandas as pd import numpy as np # 转换时间格式并排序 test_swaps['BLOCK_TIMESTAMP'] = pd.to_datetime(test_swaps['BLOCK_TIMESTAMP']) test_actions['BLOCK_TIMESTAMP'] = pd.to_datetime(test_actions['BLOCK_TIMESTAMP']) test_actions = test_actions.sort_values(['NF_TOKEN_ID', 'POOL_ADDRESS', 'BLOCK_TIMESTAMP']) # 生成每个头寸的时间区间 def process_positions(group): # 生成end_time:下一条记录的时间,最后一条设为最大交易时间(或无穷大) group['END_TIMESTAMP'] = group['BLOCK_TIMESTAMP'].shift(-1) # 获取当前池的最晚交易时间 max_swap_time = test_swaps[test_swaps['POOL_ADDRESS'] == group['POOL_ADDRESS'].iloc[0]]['BLOCK_TIMESTAMP'].max() # 最后一条记录的end_time设为max_swap_time(如果流动性非0则持续活跃) last_idx = group.index[-1] if group.loc[last_idx, 'LIQUIDITY_cumsum'] > 0: group.loc[last_idx, 'END_TIMESTAMP'] = max_swap_time else: # 注销记录的end_time不需要,因为此时流动性为0,不参与分配 group = group.iloc[:-1] return group # 按NFT和池分组处理,得到活跃头寸的时间区间、流动性和价格范围 active_positions = test_actions.groupby(['NF_TOKEN_ID', 'POOL_ADDRESS'], group_keys=False).apply(process_positions) active_positions = active_positions[active_positions['LIQUIDITY_cumsum'] > 0].reset_index(drop=True)
步骤2:按交易池拆分处理,减少关联范围
由于只有同池的交易和头寸才需要匹配,按POOL_ADDRESS拆分数据,分池处理可大幅降低计算量:
# 按交易池分组处理结果 result_list = [] for pool_addr, swaps_group in test_swaps.groupby('POOL_ADDRESS'): # 获取当前池的活跃头寸 pool_positions = active_positions[active_positions['POOL_ADDRESS'] == pool_addr] if pool_positions.empty: continue # 同池内交易与头寸关联,筛选符合条件的记录 merged = swaps_group.merge(pool_positions, on='POOL_ADDRESS', suffixes=('_swap', '_pos')) # 筛选价格在头寸区间内的记录 price_mask = (merged['PRICE_0_1'] >= merged['PRICE_LOWER_0_1']) & (merged['PRICE_0_1'] <= merged['PRICE_UPPER_0_1']) # 筛选交易时间在头寸活跃区间内的记录 time_mask = (merged['BLOCK_TIMESTAMP_swap'] >= merged['BLOCK_TIMESTAMP_pos']) & (merged['BLOCK_TIMESTAMP_swap'] <= merged['END_TIMESTAMP']) valid_records = merged[price_mask & time_mask] # 计算该头寸应得的手续费份额 valid_records['FEES_SHARE'] = (valid_records['LIQUIDITY_cumsum'] / valid_records['LIQUIDITY']) * valid_records['FEES'] # 按NFT和池分组累计手续费 pool_result = valid_records.groupby(['NF_TOKEN_ID', 'POOL_ADDRESS'])['FEES_SHARE'].sum().reset_index(name='FEES_ACCRUED') result_list.append(pool_result) # 合并所有池的结果 final_output = pd.concat(result_list, ignore_index=True)
进一步优化方向
如果同池内交易和头寸数量仍极大(比如十万级以上),可采用以下方法进一步提速:
- 用**区间索引(IntervalIndex)**处理价格匹配:将头寸的价格区间转为IntervalIndex,用
get_indexer快速找到交易价格对应的头寸; - 用merge_asof处理时间匹配:对交易按时间排序,头寸按激活时间排序,用merge_asof找到交易时间对应的活跃头寸;
- 利用Dask或PySpark进行分布式计算,处理超大规模数据。
优化效果说明
- 预处理头寸数据将活跃状态的判断从循环内重复查询转为提前计算,避免冗余操作;
- 按交易池拆分处理,将全局O(MN)复杂度降为各池O(mn)之和,大幅降低计算量;
- 矢量化操作替代双重循环,利用pandas底层C实现的优化提升运行效率。
内容的提问来源于stack exchange,提问作者Olive
相关产品推荐
相关产品推荐

