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

如何将Polars/PySpark DataFrame每行转为字典传入UDF处理?

将DataFrame每行转为字典传入UDF处理的实现方案

可以通过Polars和PySpark DataFrame实现该需求,具体方案如下:

Polars 实现方法

Polars提供了多种方式将行转为字典并传入自定义函数:

方法1:结合struct与apply(保留DataFrame格式)

import polars as pl

# 自定义处理函数
def process_row(row_dict):
    # 替换为你的业务逻辑
    return f"{row_dict['name']} ({row_dict['age']}) from {row_dict['city']}"

# 模拟目标DataFrame结构
df = pl.DataFrame({
    "name": ["Alice", "Bob", "Charlie"],
    "age": [25, 30, 35],
    "city": ["New York", "London", "Paris"]
})

# 将每行转为结构体后转字典,传入UDF处理
result_df = df.with_columns(
    pl.struct(df.columns)
      .apply(lambda row: process_row(row.to_dict()))
      .alias("processed_output")
)

print(result_df)

方法2:直接迭代行(适合逐行批量处理)

# 直接获取每行的字典形式并处理
for row_dict in df.rows(named=True):
    processed_data = process_row(row_dict)
    # 可添加存储、打印等后续逻辑
    print(processed_data)

PySpark 实现方法

PySpark可以通过RDD转换或UDF结合结构体的方式实现:

方法1:RDD映射(灵活处理单行)

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("RowToDict").getOrCreate()

# 自定义处理函数
def process_row(row_dict):
    return f"{row_dict['name']} ({row_dict['age']}) from {row_dict['city']}"

# 模拟目标DataFrame结构
df = spark.createDataFrame([
    ("Alice", 25, "New York"),
    ("Bob", 30, "London"),
    ("Charlie", 35, "Paris")
], ["name", "age", "city"])

# 将Row转为字典后传入UDF
processed_results = df.rdd.map(lambda row: process_row(row.asDict())).collect()

print(processed_results)

方法2:UDF结合struct(保留DataFrame格式)

from pyspark.sql.functions import udf, struct
from pyspark.sql.types import StringType

# 注册UDF
process_udf = udf(lambda row: process_row(row.asDict()), StringType())

# 生成结果列
result_df = df.withColumn(
    "processed_output",
    process_udf(struct(df.columns))
)

result_df.show()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 04:35:23