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

咨询:用BigQuery SQL或Python获取各客户各列最新非空值

BigQuery SQL 解决方案

核心思路

利用BigQuery的LAST_VALUE窗口函数,对每个客户的每一列按RecDate排序后取最新非空值;针对数百列的场景,通过动态SQL自动生成所有列的处理逻辑,避免手动编写重复代码。

具体实现

  1. 静态列验证示例(适合少量列测试)
    如果先验证少量列(比如col1、col2),可直接编写如下SQL:
WITH ranked_data AS (
  SELECT
    Customer,
    LAST_VALUE(col1 IGNORE NULLS) OVER (PARTITION BY Customer ORDER BY RecDate ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING) AS latest_col1,
    LAST_VALUE(col2 IGNORE NULLS) OVER (PARTITION BY Customer ORDER BY RecDate ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING) AS latest_col2,
    -- 其他列按此格式类推
    RecDate
  FROM your_dataset.your_table
)
SELECT DISTINCT
  Customer,
  latest_col1,
  latest_col2
  -- 对应添加其他列
FROM ranked_data;
  1. 动态SQL(适配数百列场景)
    针对大量列,通过查询元数据自动生成列处理逻辑:
DECLARE columns_str STRING;

-- 从元数据中提取除Customer、RecDate外的所有列
SET columns_str = (
  SELECT STRING_AGG(
    CONCAT(
      'LAST_VALUE(', column_name, ' IGNORE NULLS) OVER (PARTITION BY Customer ORDER BY RecDate ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING) AS latest_', column_name
    ),
    ', '
  )
  FROM your_dataset.INFORMATION_SCHEMA.COLUMNS
  WHERE table_name = 'your_table'
    AND column_name NOT IN ('Customer', 'RecDate')
);

-- 拼接并执行完整查询
EXECUTE IMMEDIATE CONCAT(
  'WITH ranked_data AS (
     SELECT
       Customer,
       ', columns_str, ',
       RecDate
     FROM your_dataset.your_table
   )
   SELECT DISTINCT
     Customer,
     ', REPLACE(columns_str, 'LAST_VALUE(', 'latest_'), '
   FROM ranked_data;'
);

替换your_dataset和your_table为实际数据集与表名;若需排除其他列,在column_name NOT IN中添加对应列名。


Python 解决方案(适配超大规模数据)

针对数百万客户的超大数据量,可结合BigQuery客户端或Apache Beam进行分布式处理:

方案1:BigQuery客户端分块处理

from google.cloud import bigquery

client = bigquery.Client()

# 获取目标表的所有列名(排除Customer、RecDate)
table_ref = client.dataset("your_dataset").table("your_table")
table = client.get_table(table_ref)
columns = [field.name for field in table.schema if field.name not in ("Customer", "RecDate")]

# 构造动态查询语句
columns_str = ", ".join(
    f"LAST_VALUE({col} IGNORE NULLS) OVER (PARTITION BY Customer ORDER BY RecDate) AS latest_{col}"
    for col in columns
)
query = f"""
WITH ranked_data AS (
  SELECT
    Customer,
    {columns_str},
    RecDate
  FROM `your_dataset.your_table`
)
SELECT DISTINCT Customer, {", ".join(f"latest_{col}" for col in columns)} FROM ranked_data
"""

# 将结果写入新表(自动分块处理,避免内存溢出)
destination_table = client.dataset("your_dataset").table("latest_customer_data")
client.load_table_from_query(
    query, destination_table, write_disposition="WRITE_TRUNCATE"
).result()

方案2:Apache Beam分布式处理(适配千万级以上客户)

通过Dataflow进行分布式计算,处理超大规模数据更高效:

import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions

class LatestNonEmptyFn(beam.DoFn):
    def process(self, element):
        customer, rows = element
        # 按RecDate倒序排序,优先取最新行
        sorted_rows = sorted(rows, key=lambda x: x["RecDate"], reverse=True)
        latest_values = {"Customer": customer}
        # 遍历每一列,取第一个非空值
        for col in [k for k in sorted_rows[0].keys() if k not in ("Customer", "RecDate")]:
            for row in sorted_rows:
                if row[col] is not None:
                    latest_values[f"latest_{col}"] = row[col]
                    break
        yield latest_values

def run():
    options = PipelineOptions(
        runner='DataflowRunner',
        project='your-gcp-project',
        region='us-central1',
        temp_location='gs://your-bucket/temp',
        staging_location='gs://your-bucket/staging'
    )

    with beam.Pipeline(options=options) as p:
        (p
         | '读取BigQuery数据' >> beam.io.ReadFromBigQuery(
             query='SELECT * FROM `your_dataset.your_table`',
             use_standard_sql=True
         )
         | '按客户分组' >> beam.GroupBy(lambda x: x["Customer"])
         | '提取最新非空值' >> beam.ParDo(LatestNonEmptyFn())
         | '写入结果表' >> beam.io.WriteToBigQuery(
             table='your_dataset.latest_customer_data',
             write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE,
             create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED
         )
        )

if __name__ == '__main__':
    run()

替换代码中的项目ID、数据集、表名及GCS存储桶路径为实际信息。


内容的提问来源于stack exchange,提问作者Dale Matt

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 02:20:12