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

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;

关键提醒

  1. Palantir Workbook中的数据集都是Spark分布式DataFrame,不要用Pandas的单机遍历方法(比如iterrows、for循环逐行处理),效率极低甚至报错;
  2. 优先用Spark的分组、窗口函数实现逻辑,这是Spark的原生优化方式,适配分布式场景;
  3. Python脚本最终要返回Spark DataFrame,不能返回Pandas DataFrame,否则会和Workbook的数据集体系不兼容。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 06:45:21