咨询:用BigQuery SQL或Python获取各客户各列最新非空值
BigQuery SQL 解决方案
核心思路
利用BigQuery的LAST_VALUE窗口函数,对每个客户的每一列按RecDate排序后取最新非空值;针对数百列的场景,通过动态SQL自动生成所有列的处理逻辑,避免手动编写重复代码。
具体实现
- 静态列验证示例(适合少量列测试)
如果先验证少量列(比如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;
- 动态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
相关产品推荐
相关产品推荐

