AWS Glue中PySpark任务报java.io.UncheckedIOException错误的解决咨询
AWS Glue PySpark 报错
Stream closed 解决建议 报错信息:
23/03/22 09:47:01 ERROR StreamGobbler: Exception reading inputstream of analyzer process java.io.UncheckedIOException: java.io.IOException: Stream closed
可能原因及解决步骤
1. 数据打印操作触发流异常
代码中调用的RefineResolverGroupData.print_dataframe方法如果基于show()或collect()实现,全量打印大数据量的DataFrame会占用过多输入流资源,导致流提前关闭。
- 解决:
- 调试阶段仅打印少量数据,比如用
limit(10).show()替代全量打印;生产环境直接移除这类非必要打印操作。 - 检查
print_dataframe内部实现,避免重复读取流数据的逻辑。
- 调试阶段仅打印少量数据,比如用
2. 重复引用DataFrame导致资源冲突
两次使用flatten_resolver_group_config_df执行join操作,若该DataFrame基于一次性数据源(如流数据、临时文件),重复引用会导致流被提前耗尽。
- 解决:
- 获取该DataFrame后立即缓存:
flatten_resolver_group_config_df.cache(),任务完成后释放资源:flatten_resolver_group_config_df.unpersist()。 - 或者将其转为临时视图,后续通过视图执行join:
flatten_resolver_group_config_df.createOrReplaceTempView("flatten_resolver_config")
- 获取该DataFrame后立即缓存:
3. Join操作存在列歧义或类型不匹配
两次join操作中列名未明确区分,或resolverGroup与配置表中对应列的数据类型不一致,会导致Spark分析器处理时出现流异常。
- 解决:
- 给join的表添加别名,明确指定关联列:
dfWithNonSupported = refine_resolver_group_ishspRelated_df_splitted_view.join( flatten_resolver_group_config_df.drop("nonhspSupportedResolverGroups").alias("hsp_cfg"), col("resolverGroup") == col("hsp_cfg.hspSupportedResolverGroups"), "left" ) - 统一关联列的数据类型,比如强制转为字符串类型:
col("resolverGroup").cast(StringType()) == col("hsp_cfg.hspSupportedResolverGroups").cast(StringType())
- 给join的表添加别名,明确指定关联列:
4. Glue作业资源配置不足
Executor内存、核心数配置过低,处理大数据量时资源耗尽,触发流关闭异常。
- 解决:
- 在Glue作业配置中,提升Executor内存(比如从5GB调整到10GB)和核心数。
- 启用动态分配,让Glue根据任务负载自动调整资源。
5. S3写入逻辑的流管理问题
publish_data_to_s3若使用自定义流操作而非Spark原生write API,可能存在未正确关闭输出流的情况。
- 解决:
- 优先使用Spark原生写入方法,比如
union_dataset_df.write.parquet("s3://your-bucket/path"),Spark会自动管理流生命周期。 - 若必须用自定义逻辑,确保所有流资源在写入完成后显式关闭。
- 优先使用Spark原生写入方法,比如
优化后的代码片段示例
def classify_hsp_supported_and_hsp_non_supported(self, refine_resolver_group_ishspRelated_df, refine_resolver_group_ishspNotRelated_df): refine_resolver_group_ishspRelated_df_splitted_view = refine_resolver_group_ishspRelated_df.select( "manager", explode(split(col("resolverGroup"), "[|]")).alias("resolverGroup") ) # 调试用打印,生产环境建议移除或限制行数 RefineResolverGroupData.print_dataframe(refine_resolver_group_ishspRelated_df_splitted_view.limit(10), "refine_resolver_group_ishspRelated_df_splitted_view") flatten_resolver_group_config_df = self.flatten_resolver_group_config_dataframe() # 缓存配置表,避免重复读取 flatten_resolver_group_config_df.cache() dfWithNonSupported = refine_resolver_group_ishspRelated_df_splitted_view.join( flatten_resolver_group_config_df.drop("nonhspSupportedResolverGroups").alias("hsp_cfg"), col("resolverGroup") == col("hsp_cfg.hspSupportedResolverGroups"), "left" ) dfWithSupportedAndNonSupported = dfWithNonSupported.join( flatten_resolver_group_config_df.drop("hspSupportedResolverGroups").alias("nonhsp_cfg"), col("resolverGroup") == col("nonhsp_cfg.nonhspSupportedResolverGroups"), "left" ) refine_resolver_group_df = dfWithSupportedAndNonSupported.groupBy("Manager").agg( collect_list("resolverGroup").alias("resolverGroup"), collect_list("hsp_cfg.hspSupportedResolverGroups").alias("hspSupportedResolverGroups"), collect_list("nonhsp_cfg.nonhspSupportedResolverGroups").alias("nonhspSupportedResolverGroups"), ) refine_resolver_group_df = refine_resolver_group_df.withColumn("ishspRelated", lit("t")) refine_resolver_group_df = refine_resolver_group_df.withColumn("resolverGroup", concat_ws("|", col("resolverGroup"))) refine_resolver_group_df = refine_resolver_group_df.withColumn("hspSupportedResolverGroups", concat_ws("|", col("hspSupportedResolverGroups"))) refine_resolver_group_df = refine_resolver_group_df.withColumn("nonhspSupportedResolverGroups", concat_ws("|", col("nonhspSupportedResolverGroups"))) refine_resolver_group_df = refine_resolver_group_df.select("Manager", "resolverGroup", "ishspRelated", "hspSupportedResolverGroups", "nonhspSupportedResolverGroups") # 调试用打印,生产环境建议移除或限制行数 RefineResolverGroupData.print_dataframe(refine_resolver_group_df.limit(10), "refine_resolver_group_df") union_dataset_df = RefineResolverGroupData.dataset_union(refine_resolver_group_df, refine_resolver_group_ishspNotRelated_df) refine_data_published = self.publish_data_to_s3(union_dataset_df) if not refine_data_published: logger.error("Employee Hierarchy Job异常,请检查 resolver group 处理逻辑") raise RuntimeError("Refine Resolver Group 任务失败") # 释放缓存资源 flatten_resolver_group_config_df.unpersist() logger.info("Resolver Group 数据处理完成 {}".format(run_date))
内容的提问来源于stack exchange,提问作者frp farhan
相关产品推荐
相关产品推荐

