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

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")
      

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())
      

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会自动管理流生命周期。
    • 若必须用自定义逻辑,确保所有流资源在写入完成后显式关闭。

优化后的代码片段示例

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 06:32:36