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

如何在PySpark中创建UDF结合BeautifulSoup提取HTML列p标签文本

问题原因
  • UDF使用场景错误:PySpark中@udf装饰的是行级处理函数,仅支持接收单个行的列值(即DataFrame的Column对象或字面量)作为输入。你传入的df_java.select("Body").collect()返回的是Driver端的本地Python列表,不属于UDF支持的参数类型,因此触发类型报错。
  • UDF内部逻辑错误:遍历p标签的循环中提前写了return text,第一次取到p标签文本就会直接返回,后续的文本拼接逻辑完全不会执行,即使参数类型正确也拿不到完整结果。
正确实现方案

方案1:行级处理UDF(可复用、适合分布式场景)

该UDF用于处理每行的Body列HTML内容,输出当前行所有p标签拼接后的文本,后续可按需聚合全量结果:

from pyspark.sql.types import StringType
from pyspark.sql.functions import concat_ws, collect_list
from bs4 import BeautifulSoup

# 定义UDF,指定返回类型为字符串
@udf(returnType=StringType())
def parse_html_p(html_content):
    # 空值判断避免解析异常
    if not html_content:
        return ""
    soup = BeautifulSoup(html_content, "html.parser")
    # 直接拼接当前行所有p标签文本
    return ''.join([p.get_text() for p in soup.find_all("p")])

# 用法1:给原DataFrame新增解析后的文本列,可后续做其他处理
df_with_parsed = df_java.withColumn("parsed_content", parse_html_p("Body"))

# 用法2:获取全量所有行拼接后的完整文本
full_string = df_with_parsed.agg(
    concat_ws("", collect_list("parsed_content"))
).collect()[0][0]
print(full_string)

方案2:一次性全量处理(无需UDF)

如果你只需要一次性拿到全量拼接结果,不需要分布式复用行级解析逻辑,直接用你最初的本地遍历逻辑即可,不需要额外封装UDF,性能反而更高。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 12:54:02