You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.27 06:23:13