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

AWS Glue管道过滤操作报错及数据计数异常问题求助

AWS Glue 过滤逻辑修复方案

问题根源

  1. 字符串过滤计数不符:原过滤条件中col1 >= col2或LN_ITEM_REL_QTY > '0.00'可能因字段类型(如用字符串存储数值)导致按字典序比较,而非数值比较,结果不符合预期。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 10:53:22