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

