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

如何通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 16:50:21