如何用PySpark/Pandas从带分隔符字符串动态生成列并求和排序
用 Pandas 和 PySpark 处理向量拆分、求和与排序
原始数据
| a | a_vector |
|---|---|
| 1 | 710.83;-776.98;-10.86;2013.02;-896.28; |
| 2 | 3 ; 2 ; 1 |
需求:将a_vector按顺序拆分为动态生成的col1、col2…列,计算各新列的求和结果,最后将求和值排序后放入一列。
Pandas 实现
1. 拆分向量到动态列
先清理字符串格式,再拆分扩展为多列:
import pandas as pd # 构建原始DataFrame data = { 'a': [1, 2], 'a_vector': ['710.83;-776.98;-10.86;2013.02;-896.28;', '3 ; 2 ; 1'] } df = pd.DataFrame(data) # 清理向量字符串:去空格、去末尾分号,拆分为列表 df['clean_vector'] = df['a_vector'].str.replace(' ', '').str.rstrip(';').str.split(';') # 生成动态列名 max_col_count = df['clean_vector'].str.len().max() col_names = [f'col{i+1}' for i in range(max_col_count)] # 扩展为多列并转为数值类型 df[col_names] = pd.DataFrame(df['clean_vector'].tolist(), index=df.index).astype(float) # 删除临时列 df = df.drop('clean_vector', axis=1)
2. 计算列求和并排序
# 计算各新列的总和 col_totals = df[col_names].sum().tolist() # 对总和排序 sorted_totals = sorted(col_totals) # 将排序结果添加到原表(每行重复该结果) df['sorted_col_sums'] = str(sorted_totals)
拆分后的中间结果示例:
| a | a_vector | col1 | col2 | col3 | col4 | col5 |
|---|---|---|---|---|---|---|
| 1 | 710.83;-776.98;-10.86;2013.02;-896.28; | 710.83 | -776.98 | -10.86 | 2013.02 | -896.28 |
| 2 | 3 ; 2 ; 1 | 3.0 | 2.0 | 1.0 | NaN | NaN |
PySpark 实现
1. 拆分向量到动态列
利用Spark的字符串处理和数组函数完成拆分:
from pyspark.sql import SparkSession from pyspark.sql.functions import split, col, size, element_at, regexp_replace, sum as spark_sum, array_sort, array # 初始化Spark会话 spark = SparkSession.builder.appName("vector_processing").getOrCreate() # 构建原始DataFrame data = [ (1, '710.83;-776.98;-10.86;2013.02;-896.28;'), (2, '3 ; 2 ; 1') ] df = spark.createDataFrame(data, schema=['a', 'a_vector']) # 清理向量字符串并转为数值数组 df = df.withColumn( "vector_array", split( regexp_replace(regexp_replace(col("a_vector"), " ", ""), ";$", ""), ";" ).cast("array<double>") ) # 获取最大向量长度,生成动态列名 max_col_count = df.select(size(col("vector_array"))).agg({"size(vector_array)": "max"}).collect()[0][0] col_names = [f'col{i+1}' for i in range(max_col_count)] # 动态生成拆分后的列 for idx in range(max_col_count): df = df.withColumn(col_names[idx], element_at(col("vector_array"), idx+1)) # 删除临时列 df = df.drop("vector_array")
2. 计算列求和并排序
# 计算各新列的总和 sum_exprs = [spark_sum(col(c)).alias(f'sum_{c}') for c in col_names] sum_result = df.agg(*sum_exprs) # 将求和结果转为数组并排序 sorted_sums_df = sum_result.select( array_sort(array(*[col(f'sum_{c}') for c in col_names])).alias('sorted_col_sums') ) # 将排序结果关联到原表 final_df = df.crossJoin(sorted_sums_df)
内容的提问来源于stack exchange,提问作者lunbox
相关产品推荐
相关产品推荐

