PySpark Pandas UDF实现PII掩码时遇输出行数超输入错误求助
PySpark Pandas UDF 输出行数多于输入行数的解决方案
问题根源
你的代码错误出在将整个pd.Series直接转为字符串的操作上:
strtext = str(text)中,text是pd.Series类型,直接转字符串会把整个序列(包含索引、所有元素)拼接成一个单一字符串,而非逐行处理每条通话记录。- 后续yield这个单一字符串时,Spark无法匹配输入的行数,触发
AssertionError: Pandas SCALAR_ITER UDF outputted more rows than input rows错误。
修复后的代码
import re import pandas as pd from pyspark.sql.functions import pandas_udf from typing import Iterator, Tuple @pandas_udf("string") def pu_mask_all_pii(iterator: Iterator[Tuple[pd.Series, pd.Series]]) -> Iterator[pd.Series]: for text_series, pii_lists_series in iterator: # 定义单条文本的PII掩码处理逻辑 def mask_single_entry(text, pii_list): # 按PII长度倒序排序,优先替换长匹配项,避免短PII截断长PII sorted_pii = sorted(pii_list, key=len, reverse=True) processed_text = str(text) for pii in sorted_pii: if len(pii) > 1: mask = len(pii) * 'X' # 直接处理字符串,无需encode避免额外转义 processed_text = re.sub(re.escape(pii), mask, processed_text, flags=re.IGNORECASE) return processed_text # 逐行配对处理文本与对应PII数组,生成结果Series masked_results = pd.Series([ mask_single_entry(txt, pii_list) for txt, pii_list in zip(text_series, pii_lists_series) ]) yield masked_results
关键修复点
- 逐行配对处理:通过
zip(text_series, pii_lists_series)将每条通话记录和对应的PII数组一一绑定,确保输入输出行数严格一致。 - 单条数据独立处理:将掩码逻辑封装为单行处理函数,避免批量操作导致的序列转字符串错误。
- 移除冗余编码操作:原代码中的
strtext.encode()会引入字节串标记(b''),直接处理字符串即可保证结果正确。 - 保留PII排序逻辑:继续按长度倒序排序PII,避免短匹配项先替换导致长PII无法被识别的问题。
内容的提问来源于stack exchange,提问作者Mohan Rayapuvari
相关产品推荐
相关产品推荐

