Spark 3.0.1下pydeequ回调服务器无法自动关闭,有无替代方案?
解决Pydeequ回调服务器端口残留问题
针对你遇到的调用VerificationResult.checkResultsAsDataFrame后端口25334被残留回调服务器占用的问题,除了显式关闭JavaGateway外,还可以尝试以下几种方案:
1. 避免使用Python Lambda表达式(从根源避免回调服务器启动)
Pydeequ启动回调服务器的核心原因是你在约束中使用了Python Lambda函数(比如lambda x: x < 3),这类逻辑需要Python端执行,必须通过py4j回调服务器完成。如果改用Pydeequ内置的Java侧约束逻辑,就不会触发回调服务器的启动,自然不会有端口残留问题。
比如将Lambda约束替换为Pydeequ提供的内置条件:
# 原Lambda写法 check.hasSize(lambda x: x < 3) # 替换为内置约束(根据Pydeequ版本调整) check.hasSizeLessThan(3) # 或使用预定义Condition类 from pydeequ.checks import Condition check.hasSize(Condition.isLessThan(3))
如果必须用自定义判断,也可以将逻辑迁移到Spark SQL的Java UDF中,避免Python回调依赖。
2. 直接关闭Spark关联的回调服务器(更简洁的代码写法)
不需要重新创建JavaGateway实例,直接调用Spark已有的Gateway对象关闭回调服务器即可:
checkResult_df.show(truncate=False) # 关闭回调服务器 spark.sparkContext._gateway.shutdown_callback_server() # 停止Spark Session spark.stop()
这种方式和你原有的方案效果一致,但省去了额外创建Gateway实例的步骤。
3. 将回调服务器配置为守护进程
在Spark Session初始化后,将回调服务器设置为守护进程,主进程结束时会自动回收该进程:
spark = (SparkSession .builder .config("spark.jars.packages", pydeequ.deequ_maven_coord) .config("spark.jars.excludes", pydeequ.f2j_maven_coord) .getOrCreate()) # 设置回调服务器为守护进程 spark.sparkContext._gateway.py4j_callback_server.setDaemon(True)
这种方式无需额外关闭代码,主进程结束后回调服务器会自动终止,部分环境下需验证守护进程回收机制的有效性。
4. EMR集群层面的收尾处理(备选方案)
如果代码层面调整无效,可以在EMR作业的收尾脚本中添加端口清理逻辑:
# 收尾脚本示例:清理占用25334端口的进程 lsof -i :25334 | grep LISTEN | awk '{print $2}' | xargs kill -9
这属于集群运维层面的 workaround,优先推荐代码层面的解决方案。
内容的提问来源于stack exchange,提问作者dataviews
相关产品推荐
相关产品推荐

