AWS Glue管道过滤操作报错及数据计数异常问题求助
AWS Glue 过滤逻辑修复方案
问题根源
- 字符串过滤计数不符:原过滤条件中
col1 >= col2或LN_ITEM_REL_QTY > '0.00'可能因字段类型(如用字符串存储数值)导致按字典序比较,而非数值比较,结果不符合预期。 - Filter.apply报错:
glueContext.read.jdbc()加载的是Spark DataFrame,而非Glue DynamicFrame,Filter.apply是DynamicFrame专属API,直接调用会触发参数不匹配错误。
修复方案
方案一:基于Spark DataFrame直接修正过滤逻辑
推荐使用类型安全的Column对象操作,避免字符串表达式的类型陷阱:
from pyspark.sql.functions import col # 加载数据(原代码不变) df = glueContext.read.format("jdbc")\ .option("driver", jdbc_driver_name)\ .option("url", db_url)\ .option("query", query)\ .option("user", db_username)\ .option("password", db_password)\ .load() # 前三个过滤逻辑保持不变 filtered_df0 = df.filter("ORDR_DOC_TYPE='0005'") filtered_df1 = df.filter("ORDR_DOC_TYPE='0001'") filtered_df2 = df.filter("ORDR_DOC_TYPE='0003'") # 修正第四个过滤:使用Column对象做类型安全比较 filtered_df3 = df.filter( (col("ORDR_STTS_CD") == "9000") & (col("LN_ITEM_REL_QTY") > 0.00) & (col("col1") >= col("col2")) )
如果字段是字符串存储的数值,需先转换类型再比较:
filtered_df3 = df.filter( (col("ORDR_STTS_CD") == "9000") & (col("LN_ITEM_REL_QTY").cast("double") > 0.00) & (col("col1").cast("double") >= col("col2").cast("double")) )
方案二:转成Glue DynamicFrame使用Filter.apply
若要使用Glue原生Transform API,需先转换数据结构:
from awsglue.dynamicframe import DynamicFrame from awsglue.transforms import Filter # 加载数据(原代码不变) df = glueContext.read.format("jdbc")\ .option("driver", jdbc_driver_name)\ .option("url", db_url)\ .option("query", query)\ .option("user", db_username)\ .option("password", db_password)\ .load() # 将Spark DataFrame转为Glue DynamicFrame dyf = DynamicFrame.fromDF(df, glueContext, "source_dyf") # 前三个过滤转成DynamicFrame操作(可选,也可保持原DataFrame逻辑) filtered_dyf0 = Filter.apply(frame=dyf, f=lambda x: x["ORDR_DOC_TYPE"] == "0005") filtered_dyf1 = Filter.apply(frame=dyf, f=lambda x: x["ORDR_DOC_TYPE"] == "0001") filtered_dyf2 = Filter.apply(frame=dyf, f=lambda x: x["ORDR_DOC_TYPE"] == "0003") # 修正第四个过滤:转换数值类型后再比较 filtered_dyf3 = Filter.apply( frame=dyf, f=lambda x: (x["ORDR_STTS_CD"] == "9000") and (float(x["LN_ITEM_REL_QTY"]) > 0.00) and (float(x["col1"]) >= float(x["col2"])) ) # 如需转回DataFrame进行后续处理 filtered_df3 = filtered_dyf3.toDF()
关键注意点
- 避免对数值类型字段加引号做字符串比较(如
> '0.00'),会导致非预期的字典序比较结果。 - Spark DataFrame和Glue DynamicFrame的API不通用,调用前需确认当前数据结构类型。
内容的提问来源于stack exchange,提问作者jake_jjc
相关产品推荐
相关产品推荐

