如何通过PySpark/Spark SQL提取特定格式数值并生成对应新列?
针对你提出的提取特定格式数值并生成独立新列的需求,下面分别给出PySpark和Spark SQL的实现方案:
使用PySpark实现
步骤1:定义匹配正则
需要同时匹配两种目标格式:
- 以
$开头的金额:匹配规则为\$\d+(,\d{3})*\.\d{2},覆盖带千分位逗号的场景(如$4,440.00) - 以
%结尾的百分比:匹配规则为\d+%,覆盖整数百分比场景(如25%)
将两个规则用|合并,实现一次性提取所有符合条件的数值:\$\d+(,\d{3})*\.\d{2}|\d+%
步骤2:提取、转换并拆分数值
通过正则提取数组、清洗格式、动态生成新列三个步骤完成需求:
from pyspark.sql import SparkSession from pyspark.sql.functions import regexp_extract_all, col, udf from pyspark.sql.types import ArrayType, DoubleType # 初始化Spark会话 spark = SparkSession.builder.appName("ExtractTargetValues").getOrCreate() # 模拟测试数据 test_data = [ ("$1620.00 25%",), ("$4,440.00 100% 50%",), ("无目标数值的行",), ("$999.99 30% $1,234.56",) ] df = spark.createDataFrame(test_data, ["original_col"]) # 定义清洗数值的UDF:去除格式符号并转为数字 def clean_num(val): if val.startswith("$"): return float(val.replace("$", "").replace(",", "")) elif val.endswith("%"): return float(val.replace("%", "")) return None clean_num_udf = udf(clean_num, DoubleType()) # 1. 提取所有匹配的数值到数组列 df = df.withColumn("extracted_str", regexp_extract_all(col("original_col"), r"\$\d+(,\d{3})*\.\d{2}|\d+%", 0)) # 2. 清洗数组内的每个元素,转为数字类型 df = df.withColumn("numeric_arr", clean_num_udf(col("extracted_str"))) # 3. 获取数组最大长度,动态生成对应数量的新列 max_col_count = df.selectExpr("max(size(numeric_arr))").first()[0] for idx in range(max_col_count): df = df.withColumn(f"value_{idx+1}", col("numeric_arr")[idx]) # 查看最终结果 df.show(truncate=False)
关键说明
regexp_extract_all会将每行所有匹配的数值提取为一个数组- UDF负责处理格式转换,兼容带千分位的金额和百分比
- 动态生成列的逻辑保证了不管每行有多少个目标数值,都能生成对应独立列
使用Spark SQL实现
步骤1:注册临时视图并编写SQL逻辑
通过SQL的内置函数完成提取、转换,再拆分数组为独立列:
# 沿用上述DataFrame,注册为临时视图 df.createOrReplaceTempView("target_data") # 执行SQL查询 spark.sql(""" WITH step1_extract AS ( SELECT original_col, regexp_extract_all(original_col, '\\$\\d+(,\\d{3})*\\.\\d{2}|\\d+%', 0) AS extracted_str_arr FROM target_data ), step2_clean AS ( SELECT original_col, TRANSFORM( extracted_str_arr, val -> CASE WHEN val LIKE '$%' THEN CAST(regexp_replace(val, '[$,]', '') AS DOUBLE) WHEN val LIKE '%' THEN CAST(regexp_replace(val, '%', '') AS DOUBLE) ELSE NULL END ) AS numeric_arr FROM step1_extract ), step3_get_max_len AS ( SELECT max(size(numeric_arr)) AS max_cols FROM step2_clean ) SELECT original_col, numeric_arr, numeric_arr[0] AS value_1, numeric_arr[1] AS value_2, numeric_arr[2] AS value_3 FROM step2_clean, step3_get_max_len """).show(truncate=False)
关键说明
TRANSFORM函数替代PySpark的UDF,直接在SQL中遍历数组并清洗数值- 若需要完全动态适配任意数量的列,可以通过Spark元数据API获取
max_cols值,再动态拼接SQL语句生成对应列
内容的提问来源于stack exchange,提问作者Anonymous
相关产品推荐
相关产品推荐

