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

AWS Glue作业中使用pydeequ校验数据时无法终止的问题求助

问题:AWS Glue作业调用PyDeequ部分Check方法后无法结束,LogPusher反复上传日志

我按照AWS大数据博客的步骤在Glue Studio创建作业,用PyDeequ做数据校验。PyDeequ能正常运行,但调用部分Check方法后,所有处理完成作业却一直处于运行状态。查看执行日志发现,LogPusher进程周期性输出类似glue.LogPusher (Logging.scala:logInfo(57)): uploading /tmp/spark-event-logs/ to s3://bucket/sparkHistoryLogs/的日志,反复尝试将事件日志上传至S3。

可复现的Glue作业代码:

# 作业语言:Python
# Glue版本:3.0
# Deequ版本:deequ-2.0.1-spark-3.2

from awsglue.context import GlueContext
from pyspark.context import SparkContext
from awsglue.dynamicframe import DynamicFrame
from pydeequ.analyzers import (
    AnalysisRunner,
    AnalyzerContext,
    Completeness,
    Maximum,
    MaxLength,
    Minimum,
    MinLength,
    Size,
    UniqueValueRatio,
)
from pydeequ.checks import Check, CheckLevel
from pydeequ.verification import VerificationResult, VerificationSuite

glue_context = GlueContext(SparkContext.getOrCreate())
spark = glue_context.spark_session

dyf = glue_context.create_dynamic_frame.from_options(
    format_options={"quoteChar": '"', "withHeader": True, "separator": ","},
    connection_type="s3",
    format="csv",
    connection_options={
        "paths": [f"s3://bucket/test.csv"],
        "recurse": True,
    }
)

df = dyf.toDF()
df.show()

# 输出:
# +-------+-------+
# |column1|column2|
# +-------+-------+
# |      a|      1|
# |      b|      2|
# |      c|      3|
# +-------+-------+

runner = AnalysisRunner(spark).onData(df)

runner.addAnalyzer(Size())
runner.addAnalyzer(Completeness("column1"))
runner.addAnalyzer(Completeness("column2"))
runner.addAnalyzer(UniqueValueRatio(["column1"]))
runner.addAnalyzer(UniqueValueRatio(["column2"]))
runner.addAnalyzer(MinLength("column2"))
runner.addAnalyzer(MaxLength("column2"))
runner.addAnalyzer(Minimum("column2"))
runner.addAnalyzer(Maximum("column2"))

result = runner.run()
result_df = AnalyzerContext.successMetricsAsDataFrame(spark, result)
result_df.show(truncate=False)

# 输出:
# +-------+--------+----------------+-----+
# |entity |instance|name            |value|
# +-------+--------+----------------+-----+
# |Column |column2 |UniqueValueRatio|1.0  |
# |Dataset|*       |Size            |3.0  |
# |Column |column1 |Completeness    |1.0  |
# |Column |column1 |UniqueValueRatio|1.0  |
# |Column |column2 |Completeness    |1.0  |
# |Column |column2 |MinLength       |1.0  |
# |Column |column2 |MaxLength       |1.0  |
# +-------+--------+----------------+-----+

check_1 = Check(spark, CheckLevel.Warning, "isComplete").isComplete("column1")
result_1 = VerificationSuite(spark).onData(df).addCheck(check_1).run()
result_df_1 = VerificationResult.checkResultsAsDataFrame(spark, result_1)
result_df_1.show(truncate=False)

# 输出:
# +----------+-----------+------------+--------------------------------------------------+-----------------+------------------+
# |check     |check_level|check_status|constraint                                        |constraint_status|constraint_message|
# +----------+-----------+------------+--------------------------------------------------+-----------------+------------------+
# |isComplete|Warning    |Success     |CompletenessConstraint(Completeness(column1,None))|Success          |                  |
# +----------+-----------+------------+--------------------------------------------------+-----------------+------------------+

# 至此,作业可成功完成。

check_2 = Check(spark, CheckLevel.Warning, "hasMinLength").hasMinLength("column1",lambda x: x >= 1)
result_2 = VerificationSuite(spark).onData(df).addCheck(check_2).run()
result_df_2 = VerificationResult.checkResultsAsDataFrame(spark, result_2)
result_df_2.show(truncate=False)

# 输出:
# +------------+-----------+------------+--------------------------------------------+-----------------+------------------+
# |check       |check_level|check_status|constraint                                  |constraint_status|constraint_message|
# +------------+-----------+------------+--------------------------------------------+-----------------+------------------+
# |hasMinLength|Warning    |Success     |MinLengthConstraint(MinLength(column1,None))|Success          |                  |
# +------------+-----------+------------+--------------------------------------------+-----------------+------------------+

# 执行上述流程后,结果正常显示,但作业永远无法结束。

原因分析

  • PyDeequ的lambda序列化兼容性问题:在hasMinLength等Check方法中使用Python lambda表达式作为约束条件时,PySpark无法将lambda正确序列化并传递给Scala后端的Deequ组件。这会导致Spark作业中残留未关闭的线程或资源,阻止Glue作业正常终止。
  • LogPusher的触发逻辑:AWS Glue的LogPusher进程会持续尝试上传Spark事件日志,直到作业完全终止。由于资源未释放,作业无法进入结束状态,LogPusher会反复执行上传操作。

解决方法

1. 替换lambda为Deequ内置约束方法

对于简单数值约束,直接使用Deequ提供的参数化方法,避免Python lambda:

# 替换lambda表达式为数值参数
check_2 = Check(spark, CheckLevel.Warning, "hasMinLength").hasMinLength("column1", 1)

2. 显式关闭Spark上下文(应急方案)

若必须使用复杂逻辑无法避免lambda,可在作业最后强制关闭Spark上下文释放资源:

# 在所有业务逻辑执行完成后添加
spark.stop()

注意:此方法可能导致部分日志未正常上传,但能确保作业终止。

3. 升级PyDeequ版本

检查PyDeequ最新版本,部分新版本修复了Python-Scala序列化兼容问题,可尝试升级到适配Spark 3.2的更高版本(如deequ-2.0.2-spark-3.2及以上)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 13:11:37