【问题标题】:How configure glue bookmars to work with scala code?如何配置胶水书签以使用 scala 代码?
【发布时间】:2018-01-18 05:47:35
【问题描述】:

考虑scala代码:

import com.amazonaws.services.glue.GlueContext
import com.amazonaws.services.glue.util.{GlueArgParser, Job, JsonOptions}
import org.apache.spark.SparkContext

import scala.collection.JavaConverters.mapAsJavaMapConverter

object MyGlueJob {

  def main(sysArgs: Array[String]) {
    val spark: SparkContext = SparkContext.getOrCreate()
    val glueContext: GlueContext = new GlueContext(spark)

    val args = GlueArgParser.getResolvedOptions(sysArgs, Seq("JOB_NAME").toArray)
    Job.init(args("JOB_NAME"), glueContext, args.asJava)

    val input = glueContext
      .getCatalogSource(database = "my_data_base", tableName = "my_json_gz_partition_table")
      .getDynamicFrame()

    val processed = input.applyMapping(
      Seq(
        ("id",                                        "string", "id", "string"),
        ("my_date",                                   "string", "my_date", "string")
      ))
    glueContext.getSinkWithFormat(
      connectionType = "s3",
      options = JsonOptions(Map("path" -> "s3://my_path", "partitionKeys" -> List("my_date"))),
      format = "orc", transformationContext = ""
    ).writeDynamicFrame(processed)
    Job.commit
  }
}

输入是经过 gzip 压缩的分区 json 文件,按日期列分区。一切工作 - 数据以 json 格式读取并用 orc 编写。

但是当尝试使用相同的数据运行作业时,它会再次读取它并写入重复的数据。此作业中启用了书签。调用方法 Job.initJob.commit。怎么了?

更新

我在getCatalogSourcegetSinkWithFormat 中添加了transformationContext 参数:

        val input = glueContext
      .getCatalogSource(database = "my_data_base", tableName = "my_json_gz_partition_table", transformationContext = "transformationContext1")
      .getDynamicFrame()

和:

    glueContext.getSinkWithFormat(
      connectionType = "s3",
      options = JsonOptions(Map("path" -> "s3://my_path", "partitionKeys" -> List("my_date"))),
      format = "orc", transformationContext = "transformationContext2"
    ).writeDynamicFrame(processed)

现在魔法以这种方式“起作用”:

  1. 第一次运行 - 好的
  2. 第二次运行(使用相同的数据或相同的数据和新的数据) - 失败并出现错误(稍后)

在第二次(及后续)运行后再次发生错误。 消息Skipping Partition {"my_date": "2017-10-10"} 也会出现在日志中。

ERROR ApplicationMaster: User class threw exception: org.apache.spark.sql.AnalysisException: Partition column my_date not found in schema StructType(); org.apache.spark.sql.AnalysisException: Partition column my_date not found in schema StructType();
at org.apache.spark.sql.execution.datasources.PartitioningUtils$$anonfun$partitionColumnsSchema$1$$anonfun$apply$11.apply(PartitioningUtils.scala:439)
at org.apache.spark.sql.execution.datasources.PartitioningUtils$$anonfun$partitionColumnsSchema$1$$anonfun$apply$11.apply(PartitioningUtils.scala:439)
at scala.Option.getOrElse(Option.scala:121)
at org.apache.spark.sql.execution.datasources.PartitioningUtils$$anonfun$partitionColumnsSchema$1.apply(PartitioningUtils.scala:438)
at org.apache.spark.sql.execution.datasources.PartitioningUtils$$anonfun$partitionColumnsSchema$1.apply(PartitioningUtils.scala:437)
at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234)
at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234)
at scala.collection.immutable.List.foreach(List.scala:381)
at scala.collection.TraversableLike$class.map(TraversableLike.scala:234)
at scala.collection.immutable.List.map(List.scala:285)
at org.apache.spark.sql.execution.datasources.PartitioningUtils$.partitionColumnsSchema(PartitioningUtils.scala:437)
at org.apache.spark.sql.execution.datasources.PartitioningUtils$.validatePartitionColumn(PartitioningUtils.scala:420)
at org.apache.spark.sql.execution.datasources.DataSource.write(DataSource.scala:443)
at org.apache.spark.sql.DataFrameWriter.save(DataFrameWriter.scala:215)
at com.amazonaws.services.glue.SparkSQLDataSink.writeDynamicFrame(DataSink.scala:123)
at MobileArcToRaw$.main(script_2018-01-18-08-14-38.scala:99)

胶水书签到底是怎么回事???哦

【问题讨论】:

    标签: scala amazon-web-services aws-glue


    【解决方案1】:

    您是否尝试将transformationContext 的值设置为源和接收器的值相同?它们目前在您上次更新时设置为不同的值。

    transformationContext = "transformationContext1"

    transformationContext = "transformationContext2"

    我也曾在使用 Glue 和书签时遇到过这个问题。我正在尝试执行类似的任务,我读入按年、月和日分区的分区 JSON 文件,每天都有新文件到达。我的作业运行转换以提取数据子集,然后沉入 S3 上的分区 Parquet 文件中。

    我使用的是 Python,所以我的 DynamicFrame 的初始实例如下所示:

    dyf = glue_context.create_dynamic_frame.from_catalog(database="dev-db", table_name="raw", transformation_ctx="raw")

    最后像这样接收到 S3:

    glue_context.write_dynamic_frame.from_options( frame=select_out, connection_type='s3', connection_options={'path': output_dir, 'partitionKeys': ['year', 'month', 'day']}, format='parquet', transformation_ctx="dev-transactions" )

    最初我运行了该作业,并且启用了书签后正确生成了 Parquet。然后我添加了新的一天数据,更新了输入表上的分区并重新运行。第二个作业将失败并出现如下错误:

    pyspark.sql.utils.AnalysisException: u"cannot resolve 'year' given input columns: [];;\n'Project ['year, 'month, 'day, 'data']

    transformation_ctx 更改为相同(在我的情况下为dev-transactions)使该过程能够在仅处理增量分区并为新分区生成 Parquet 的情况下正常工作。

    关于 Bookmarks 以及如何使用转换上下文变量的文档非常少。

    Python 文档只是说:(https://docs.aws.amazon.com/glue/latest/dg/aws-glue-api-crawler-pyspark-extensions-glue-context.html):

    transformation_ctx – 要使用的转换上下文(可选)。

    Scala 文档说 (https://docs.aws.amazon.com/glue/latest/dg/glue-etl-scala-apis-glue-gluecontext.html):

    transformationContext — 与作业书签使用的接收器关联的转换上下文。默认设置为空。

    由于文档在解释方面做得很差,我能观察到的最佳结果是,转换上下文用于在已处理的源数据和接收器数据之间形成联系,并且定义不同的上下文会阻止书签按预期工作.

    【讨论】:

      【解决方案2】:

      该作业第二次运行时,似乎没有为您的目录找到新数据

      val input = glueContext.getCatalogSource(...)
      input.count
      # Returns 0, your dynamic frame has no Schema associated to it
      # hence the `Partition column my_date not found in schema StructType()`
      

      我建议在尝试映射/写入之前检查 DynamicFrame 的大小,或者您的分区字段是否存在于 DynamicFrame input.schema.containsField("my_field") 的架构中。那时,您可以提交或不提交工作。

      此外,如果您确定新数据将在新分区上进入该目录,则可以考虑运行 Crawler 来挑选这些新分区,或者如果您不希望架构发生任何更改,则可以通过 API 创建它们。

      希望这会有所帮助。

      【讨论】:

      • input.count - 它就像扫描所有数据(如果它们存在),这意味着我扫描 DynamicFrame 2 次。 :(
      • 您可以查看架构,input.schema.containsField("my_date")
      • 首先,如果没有数据,胶水会导致工作失败,这看起来很奇怪。如果是真的,那就是严重的架构错误。其次,关于input.schema内部schema通过调用DynamicFrame内部的records()方法计算得到数据读取。我认为现在没有数据扫描就无法获得模式:(
      【解决方案3】:

      JobBookmarks 使用转换上下文来关闭给定 ETL 操作的状态(主要是源)。目前将它们放在水槽中没有任何影响。

      启用作业书签时作业失败的原因之一是因为它们只处理增量数据(新文件),如果没有新数据,脚本会像没有数据时一样运行,这可以以spark分析异常为例。

      因此,您不应在不同的 ETL 运算符之间使用相同的转换上下文。

      对于您第一次运行后的测试,请尝试将新数据复制到您的源位置并再次运行该作业,应该只处理新数据。

      【讨论】:

        猜你喜欢
        • 2020-11-22
        • 1970-01-01
        • 1970-01-01
        • 2021-10-24
        • 2011-04-09
        • 1970-01-01
        • 2017-04-20
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多