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

PySpark中过滤DataFrame特定行及解决_get_object_id错误

问题解决:Spark AttributeError: 'DataFrame' object has no attribute '_get_object_id'

需求说明

过滤掉prev_df中Wellname列值存在于input_df的Wellname列中的行。

示例数据

input_df

Wellname WellType Platform
0       E17 Producer      DWG
1       E17 Producer      DWG
2       E17 Producer      DWG
3      E20Y Producer      DWG
4      E20Y Producer      DWG
5      E20Y Producer      DWG
6      E20Y Producer      DWG
7      E20Y Producer      DWG

prev_df

Wellname WellType Platform
0       E17 Producer      CH
1       E17 Producer      CH
2       E17 Producer      CH
3       E21 Producer      DWG
4       E21 Producer      DWG
5       E21 Producer      DWG

过滤后结果

Wellname WellType Platform
0       E21 Producer      CH
1       E21 Producer      CH
2       E21 Producer      CH

错误信息

spark AttributeError: 'DataFrame' object has no attribute '_get_object_id'

原代码

from transforms.api import transform, incremental, Input, Output, configure
from pyspark.sql import types as T   
schema = T.StructType([
    T.StructField('Wellname', T.StringType()),
    T.StructField('Platform', T.StringType()),
    T.StructField('date', T.TimestampType()),
    T.StructField('oil_STB_d', T.DoubleType()),
    T.StructField('water_STB_d', T.DoubleType()),
    T.StructField('GOR_scf_STB', T.DoubleType()),
    T.StructField('WHP_psi', T.DoubleType()),
    T.StructField('BHP_psi', T.DoubleType()),
    T.StructField('predicted_T_Celcius', T.DoubleType()),
    T.StructField('Precipitated_wax_wt_percentage_with_predicted_T_minus_5_Celcius', T.DoubleType()),
    T.StructField('Rate_of_wax_precipitation_with_predicted_T_minus_5_Celcius_kg_per_d', T.DoubleType()),
    T.StructField('Precipitated_wax_wt_percentage_with_predicted_T_0_Celcius', T.DoubleType()),
    T.StructField('Rate_of_wax_precipitation_with_predicted_T_0_Celcius_kg_per_d', T.DoubleType()),
    T.StructField('Precipitated_wax_wt_percentage_with_predicted_T_plus_5_Celcius', T.DoubleType()),
    T.StructField('Rate_of_wax_precipitation_with_predicted_T_plus_5_Celcius_kg_per_d', T.DoubleType()),
    T.StructField('WAT_DegC', T.DoubleType()),
    T.StructField('Wax_Content_percentage', T.DoubleType())
])    

@configure(profile=['KUBERNETES_NO_EXECUTORS'])
@incremental(require_incremental=True, snapshot_inputs=["input_df"])
@transform(
    input_df=Input('ri.foundry.lava-catalog.dataset.3c85fd4e-ca5d-4f9e-bbab-c7d3aef1da9e'),
    output_df=Output('ri.foundry.lava-catalog.dataset.44012abe-87d6-4d38-802c-16a2ed447e07')
)
def incremental_filter(input_df, output_df):
    new_df = input_df.dataframe()

    prev_df = output_df.dataframe('previous', schema)
    #Error in the following line
    prev_df = prev_df.filter(~prev_df.Wellname.isin(new_df.select('Wellname').distinct()))
    new_df = new_df.unionByName(prev_df)
    mode = 'replace'

    output_df.set_mode(mode)
    output_df.write_dataframe(new_df)

问题原因与解决方案

原因

isin()方法要求传入值列表,而非DataFrame。直接传入new_df.select('Wellname').distinct()这个DataFrame,Spark无法解析该类型参数,从而触发_get_object_id相关的属性错误。

解决方法

方法1:提取值列表后过滤

先把input_df中去重的Wellname转为Python列表,再传入isin():

from transforms.api import transform, incremental, Input, Output, configure
from pyspark.sql import types as T   
schema = T.StructType([
    T.StructField('Wellname', T.StringType()),
    T.StructField('Platform', T.StringType()),
    T.StructField('date', T.TimestampType()),
    T.StructField('oil_STB_d', T.DoubleType()),
    T.StructField('water_STB_d', T.DoubleType()),
    T.StructField('GOR_scf_STB', T.DoubleType()),
    T.StructField('WHP_psi', T.DoubleType()),
    T.StructField('BHP_psi', T.DoubleType()),
    T.StructField('predicted_T_Celcius', T.DoubleType()),
    T.StructField('Precipitated_wax_wt_percentage_with_predicted_T_minus_5_Celcius', T.DoubleType()),
    T.StructField('Rate_of_wax_precipitation_with_predicted_T_minus_5_Celcius_kg_per_d', T.DoubleType()),
    T.StructField('Precipitated_wax_wt_percentage_with_predicted_T_0_Celcius', T.DoubleType()),
    T.StructField('Rate_of_wax_precipitation_with_predicted_T_0_Celcius_kg_per_d', T.DoubleType()),
    T.StructField('Precipitated_wax_wt_percentage_with_predicted_T_plus_5_Celcius', T.DoubleType()),
    T.StructField('Rate_of_wax_precipitation_with_predicted_T_plus_5_Celcius_kg_per_d', T.DoubleType()),
    T.StructField('WAT_DegC', T.DoubleType()),
    T.StructField('Wax_Content_percentage', T.DoubleType())
])    

@configure(profile=['KUBERNETES_NO_EXECUTORS'])
@incremental(require_incremental=True, snapshot_inputs=["input_df"])
@transform(
    input_df=Input('ri.foundry.lava-catalog.dataset.3c85fd4e-ca5d-4f9e-bbab-c7d3aef1da9e'),
    output_df=Output('ri.foundry.lava-catalog.dataset.44012abe-87d6-4d38-802c-16a2ed447e07')
)
def incremental_filter(input_df, output_df):
    new_df = input_df.dataframe()

    prev_df = output_df.dataframe('previous', schema)
    
    # 提取去重后的Wellname值列表
    exclude_wells = [row.Wellname for row in new_df.select('Wellname').distinct().collect()]
    
    # 非空列表才执行过滤,避免空列表过滤掉全部数据
    if exclude_wells:
        prev_df = prev_df.filter(~prev_df.Wellname.isin(exclude_wells))
    
    new_df = new_df.unionByName(prev_df)
    mode = 'replace'

    output_df.set_mode(mode)
    output_df.write_dataframe(new_df)

方法2:左反连接(大数据量推荐)

如果数据量较大,collect()会占用Driver内存,推荐用Spark的左反连接实现过滤,效率更高:

# 替换原过滤逻辑行
prev_df = prev_df.join(new_df.select('Wellname').distinct(), on='Wellname', how='left_anti')

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 05:50:58