基于LIFO规则的交易流水配对:BigQuery SQL实现需求
LIFO交易匹配:BigQuery SQL与Python非递归实现方案
需求说明
需对千万级客户的交易表(单客户数百条交易)按**LIFO(后进先出)**规则完成Outflow(流出)与Inflow(流入)的映射匹配,核心规则:
- 优先用后续流出匹配最新的流入记录
- 若流出金额超出当前剩余流入额度,向前追溯更早的流入
- 保留未匹配的流入记录用于后续映射
部署环境为GCP BigQuery,要求非递归实现(已有SAS版本,需转SQL/Python)
BigQuery SQL 非递归实现
前提假设
交易表结构如下:
`transactions` ( customer_id STRING, transaction_id STRING, -- 唯一交易标识 transaction_type STRING, -- 取值'INFLOW'/'OUTFLOW' amount NUMERIC, transaction_timestamp TIMESTAMP )
实现代码
WITH inflows AS ( -- 预处理流入数据:按客户分组,时间倒序排序,计算累计流入及初始剩余额度 SELECT customer_id, transaction_id AS inflow_id, amount AS inflow_amount, transaction_timestamp AS inflow_ts, SUM(amount) OVER ( PARTITION BY customer_id ORDER BY transaction_timestamp DESC ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS cumulative_inflow, amount AS remaining_inflow FROM `your-project.your-dataset.transactions` WHERE transaction_type = 'INFLOW' ), outflows AS ( -- 预处理流出数据:按客户分组,时间正序排序,计算累计流出 SELECT customer_id, transaction_id AS outflow_id, amount AS outflow_amount, transaction_timestamp AS outflow_ts, SUM(amount) OVER ( PARTITION BY customer_id ORDER BY transaction_timestamp ASC ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS cumulative_outflow FROM `your-project.your-dataset.transactions` WHERE transaction_type = 'OUTFLOW' ), matched_pairs AS ( -- 基于累计值区间匹配流入与流出 SELECT i.customer_id, i.inflow_id, o.outflow_id, i.inflow_ts, o.outflow_ts, -- 计算当前流出可匹配该流入的金额 LEAST( i.remaining_inflow, o.outflow_amount, i.cumulative_inflow - COALESCE(LAG(o.cumulative_outflow) OVER (PARTITION BY i.customer_id ORDER BY o.cumulative_outflow), 0) ) AS matched_amount, -- 更新流入剩余额度 i.remaining_inflow - LEAST( i.remaining_inflow, o.outflow_amount, i.cumulative_inflow - COALESCE(LAG(o.cumulative_outflow) OVER (PARTITION BY i.customer_id ORDER BY o.cumulative_outflow), 0) ) AS updated_remaining_inflow FROM inflows i JOIN outflows o ON i.customer_id = o.customer_id -- 匹配条件:流出累计值落在流入累计区间内 AND o.cumulative_outflow <= i.cumulative_inflow AND COALESCE(LAG(o.cumulative_outflow) OVER (PARTITION BY i.customer_id ORDER BY o.cumulative_outflow), 0) < i.cumulative_inflow ), final_inflows AS ( -- 聚合得到每个流入的最终剩余额度,包含完全未匹配的流入 SELECT customer_id, inflow_id, inflow_ts, inflow_amount, MAX(updated_remaining_inflow) AS remaining_inflow FROM matched_pairs GROUP BY customer_id, inflow_id, inflow_ts, inflow_amount UNION ALL SELECT customer_id, inflow_id, inflow_ts, inflow_amount, inflow_amount AS remaining_inflow FROM inflows WHERE inflow_id NOT IN (SELECT DISTINCT inflow_id FROM matched_pairs) ) -- 输出最终匹配结果及剩余流入 SELECT fp.customer_id, fp.inflow_id, fp.inflow_ts, fp.inflow_amount, fp.remaining_inflow, mp.outflow_id, mp.outflow_ts, mp.matched_amount FROM final_inflows fp LEFT JOIN matched_pairs mp ON fp.customer_id = mp.customer_id AND fp.inflow_id = mp.inflow_id ORDER BY fp.customer_id, fp.inflow_ts DESC, mp.outflow_ts ASC;
方案优势
- 完全基于窗口函数和区间匹配,无递归逻辑,适配BigQuery大规模并行处理能力
- 按
customer_id分区处理,避免跨客户数据干扰 - 自动保留未匹配的流入记录,满足后续映射需求
Python 实现(基于Pandas+BigQuery API)
适合需要灵活逻辑调整或分批次处理的场景,核心思路是按客户分组,逐笔匹配最新流入与流出:
import pandas as pd from google.cloud import bigquery # 初始化BigQuery客户端 client = bigquery.Client(project="your-project") # 读取交易数据:按客户分组,流入倒序排序 query = """ SELECT customer_id, transaction_id, transaction_type, amount, transaction_timestamp FROM `your-project.your-dataset.transactions` ORDER BY customer_id, transaction_timestamp DESC """ df = client.query(query).to_dataframe() def process_customer_lifo(group): """处理单个客户的LIFO匹配逻辑""" # 拆分流入(倒序)和流出(正序) inflows = group[group["transaction_type"] == "INFLOW"].copy().reset_index(drop=True) outflows = group[group["transaction_type"] == "OUTFLOW"].copy() outflows = outflows.sort_values("transaction_timestamp").reset_index(drop=True) inflows["remaining"] = inflows["amount"] matched_records = [] # 遍历每笔流出,从最新流入开始匹配 for _, outflow in outflows.iterrows(): remaining_outflow = outflow["amount"] # 筛选有剩余额度的流入,按倒序(最新在前)匹配 available_inflows = inflows[inflows["remaining"] > 0].copy() for idx, inflow in available_inflows.iterrows(): match_amount = min(remaining_outflow, inflow["remaining"]) # 记录匹配结果 matched_records.append({ "customer_id": group.name, "inflow_id": inflow["transaction_id"], "outflow_id": outflow["transaction_id"], "matched_amount": match_amount, "inflow_ts": inflow["transaction_timestamp"], "outflow_ts": outflow["transaction_timestamp"] }) # 更新流入剩余额度 inflows.loc[idx, "remaining"] -= match_amount remaining_outflow -= match_amount if remaining_outflow == 0: break # 整理剩余流入记录 remaining_inflows = inflows[inflows["remaining"] > 0].rename( columns={ "transaction_id": "inflow_id", "transaction_timestamp": "inflow_ts", "amount": "inflow_amount" } )[["customer_id", "inflow_id", "inflow_ts", "inflow_amount", "remaining"]] # 合并匹配结果与剩余流入 matched_df = pd.DataFrame(matched_records) if matched_records else pd.DataFrame(columns=[ "customer_id", "inflow_id", "outflow_id", "matched_amount", "inflow_ts", "outflow_ts" ]) final_df = pd.merge( remaining_inflows, matched_df, on=["customer_id", "inflow_id", "inflow_ts"], how="left" ) return final_df # 按客户分组处理LIFO匹配 result_df = df.groupby("customer_id").apply(process_customer_lifo).reset_index(drop=True) # 将结果写回BigQuery load_job = client.load_table_from_dataframe( result_df, "your-project.your-dataset.lifo_matched_results" ) load_job.result() # 等待写入完成
方案说明
- 按客户分组处理,避免内存过载,可结合BigQuery的分桶查询实现分批处理
- 逻辑直观,便于自定义调整匹配规则
- 适合单客户交易数据量不大的场景,若数据量过大可考虑分批次处理
内容的提问来源于stack exchange,提问作者Gautham Ganesan
相关产品推荐
相关产品推荐

