Azure Databricks中Pandas UDF迭代器输出长度不匹配问题求助
解决Azure Databricks中Scalar迭代器Pandas UDF的输出长度不匹配问题
问题原因
你的代码报错核心在于Scalar迭代器类型的Pandas UDF要求输入与输出的行数严格一致,但当前实现中,你通过raw_html.to_string(index=False)把整个输入pd.Series(包含多行数据)转换成了单个字符串,再包装成单元素pd.Series返回,导致当输入batch包含多行时,输出只有一行,触发长度不匹配的错误。
修正后的代码
import re import pandas as pd from pyspark.sql import SparkSession from pyspark.sql.functions import pandas_udf, StringType from typing import Iterator # 测试数据(添加一行后) data = [ {"inputData":"<html>Tanuj is older than Eina. Chetan is older than Tanuj. Eina is older than Chetan. If the first 2 statements are true, the 3rd statement is"}, {"inputData":"<html>Pens cost more than pencils. Pens cost less than eraser. Erasers cost more than pencils and pens. If the first two statements are true, the third statement is"}, {"inputData":"<html>If we have a tree of n nodes, how many edges will it have?"}, {"inputData":"<div>Which of the following data structures can handle updates and queries in log(n) time on an array?"}, {"inputData":"<p>What is the time complexity of binary search on a sorted array?"} ] spark = SparkSession.builder.getOrCreate() df = spark.createDataFrame(data) # 修正后的HTML清洗UDF @pandas_udf(StringType()) def clean_html(raw_htmls: Iterator[pd.Series]) -> Iterator[pd.Series]: pd.set_option('display.max_colwidth', 10000) # 预编译正则,避免循环内重复编译提升性能 clean_tag_regx = re.compile("<.*?>|&([a-z0-9]+|#0-9{1,6}|#x[0-9a-f]{1,6});") multi_space_regx = re.compile(r"\s+") for raw_html in raw_htmls: # 对Series每个元素单独执行清洗,保持行数一致 cleantext = raw_html.str.replace(clean_tag_regx, " ", regex=True) cleantext = cleantext.str.replace(multi_space_regx, " ", regex=True) cleantext = cleantext.str.strip() # 可选:去除首尾空格 yield cleantext df = df.withColumn("Question", clean_html("inputData")) display(df)
关键修正点
- 避免合并Series为单个字符串:移除
to_string()操作,改用pandas的str.replace()矢量化方法,对输入Series的每个元素单独处理,确保输出Series的行数与输入完全一致。 - 预编译正则表达式:将正则编译逻辑移到循环外部,减少重复编译的性能开销。
- 可选优化:添加
str.strip()去除文本首尾的多余空格,提升清洗后文本的整洁度。
内容的提问来源于stack exchange,提问作者Ancil Pa
相关产品推荐
相关产品推荐

