【发布时间】:2022-08-13 05:29:39
【问题描述】:
我正在使用 pydeequ 和 Spark 3.0.1 对数据执行一些约束检查。
至于用VerificationSuite测试,在调用VerificationResult.checkResultsAsDataFrame(spark, result)之后,pydeequ启动的回调服务器似乎没有自动终止。
例如,如果我在 EMR 集群上运行包含 pydeequ 的代码,则端口 25334 在 spark 应用程序关闭后似乎保持打开状态,除非我使用 spark 会话显式创建 JavaGateway 并调用 close() 方法。
from pydeequ.verification import *
from pyspark.sql import SparkSession, Row
spark = (SparkSession
.builder
.config(\"spark.jars.packages\", pydeequ.deequ_maven_coord)
.config(\"spark.jars.excludes\", pydeequ.f2j_maven_coord)
.getOrCreate())
df = spark.sparkContext.parallelize([
Row(a=\"foo\", b=1, c=5),
Row(a=\"bar\", b=2, c=6),
Row(a=\"baz\", b=None, c=None)]).toDF()
from py4j.java_gateway import JavaGateway
check = Check(spark, CheckLevel.Warning, \"Review Check\")
checkResult = VerificationSuite(spark) \\
.onData(df) \\
.addCheck(
check.hasSize(lambda x: x < 3) \\
.hasMin(\"b\", lambda x: x == 0) \\
.isComplete(\"c\") \\
.isUnique(\"a\") \\
.isContainedIn(\"a\", [\"foo\", \"bar\", \"baz\"]) \\
.isNonNegative(\"b\")) \\
.run()
checkResult_df = VerificationResult.checkResultsAsDataFrame(spark, checkResult)
checkResult_df.show(truncate=False)
a = JavaGateway(spark.sparkContext._gateway)
a.close()
如果我没有实现最后两行代码,回调服务器在端口上保持打开状态。
有没有解决的办法?
标签: apache-spark pyspark amazon-deequ pydeequ