Palantir Workbook中遍历数据集并生成新数据集的技术求助
解决方案:Palantir Workbook处理患者多行数据(Spark SQL/Python)
报错原因
你用的SourceDataSet是Spark DataFrame,不是Pandas DataFrame,所以没有iterrows()方法——这是Palantir Workbook里的核心数据类型,和你熟悉的Pandas本地DataFrame逻辑完全不同,不能用单机遍历的方式处理。
可行方案:Python(Spark原生API)
Spark是分布式计算框架,要利用它的分组、窗口函数来实现按患者遍历行的逻辑,以下是针对你的需求的示例代码:
示例1:提取患者ID(替代原iterrows逻辑)
def process_patient_data(SourceDataSet): from pyspark.sql import functions as F # 直接选择需要的字段,Spark会分布式执行,无需手动遍历 return SourceDataSet.select("patient_icn")
示例2:识别重复行(按患者+其他细节字段判断)
def process_patient_data(SourceDataSet): from pyspark.sql import functions as F # 按所有字段分组,统计重复次数 duplicate_df = SourceDataSet.groupBy(*SourceDataSet.columns) \ .agg(F.count("*").alias("duplicate_count")) # 筛选出重复的行(次数>1) return duplicate_df.filter(F.col("duplicate_count") > 1)
示例3:根据历史行判断状态(比如状态变化、最新记录)
假设数据集有patient_icn、status、record_time字段,用窗口函数实现按患者分组的历史对比:
def process_patient_data(SourceDataSet): from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义窗口:按患者分组,按记录时间倒序排序 patient_window = Window.partitionBy("patient_icn") \ .orderBy(F.col("record_time").desc()) # 添加列:标记是否为最新记录、获取上一条记录的状态 result_df = SourceDataSet.withColumn("is_latest_record", F.row_number().over(patient_window) == 1) \ .withColumn("previous_status", F.lag("status").over(patient_window)) return result_df
可行方案:Spark SQL(适合MS SQL背景)
如果你更熟悉SQL语法,Palantir Workbook支持直接写Spark SQL,用窗口函数就能实现需求:
示例1:提取患者ID
SELECT patient_icn FROM SourceDataSet;
示例2:识别重复行
SELECT *, COUNT(*) OVER (PARTITION BY patient_icn, detail_field1, detail_field2) AS duplicate_count FROM SourceDataSet WHERE duplicate_count > 1;
注:把detail_field1, detail_field2替换成你用来判断重复的具体字段
示例3:判断状态变化/最新记录
SELECT patient_icn, status, record_time, -- 标记是否为患者的最新记录 ROW_NUMBER() OVER (PARTITION BY patient_icn ORDER BY record_time DESC) AS row_rank, -- 获取患者上一条记录的状态 LAG(status) OVER (PARTITION BY patient_icn ORDER BY record_time) AS prev_status FROM SourceDataSet;
关键提醒
- Palantir Workbook中的数据集都是Spark分布式DataFrame,不要用Pandas的单机遍历方法(比如
iterrows、for循环逐行处理),效率极低甚至报错; - 优先用Spark的分组、窗口函数实现逻辑,这是Spark的原生优化方式,适配分布式场景;
- Python脚本最终要返回Spark DataFrame,不能返回Pandas DataFrame,否则会和Workbook的数据集体系不兼容。
内容的提问来源于stack exchange,提问作者WallyCode
相关产品推荐
相关产品推荐

