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
相关产品推荐
相关产品推荐

