如何将 csv/txt 文件加载到 AWS Glue 作业中

RK.*_*RK. 2 pyspark aws-glue

我对 AWS Glue 有以下 2 个说明,请您澄清一下。因为我需要在我的项目中使用胶水。

  1. 我想将 csv/txt 文件加载到 Glue 作业中进行处理。(就像我们在 Spark 中使用数据帧所做的那样)。这在胶水中可能吗?或者我们是否必须只使用 Crawler 将数据抓取到 Glue 表中并像下面一样使用它们进行进一步处理?

    empdf = glueContext.create_dynamic_frame.from_catalog(
        database="emp",
        table_name="emp_json")
    
    Run Code Online (Sandbox Code Playgroud)
  2. 下面我使用 Spark 代码将文件加载到 Glue 中,但我收到了冗长的错误日志。我们可以直接运行 Spark 或 PySpark 代码而无需对 Glue 进行任何更改吗?

    import sys
    from pyspark.context import SparkContext
    from awsglue.context import GlueContext
    
    sc = SparkContext()
    glueContext = GlueContext(sc)
    spark = glueContext.spark_session
    job = Job(glueContext)
    job.init(args['JOB_NAME'], args)
    dfnew = spark.read.option("header","true").option("delimiter", ",").csv("C:\inputs\TEST.txt")
    dfnew.show(2)
    
    Run Code Online (Sandbox Code Playgroud)

Yur*_*ruk 8

可以使用 Glue 直接从 s3 加载数据:

sourceDyf = glueContext.create_dynamic_frame_from_options(
    connection_type="s3",
    format="csv",
    connection_options={
        "paths": ["s3://bucket/folder"]
    },
    format_options={
        "withHeader": True,
        "separator": ","
    })
Run Code Online (Sandbox Code Playgroud)

你也可以只用 spark 来做到这一点(正如你已经尝试过的那样):

sourceDf = spark.read
    .option("header","true")
    .option("delimiter", ",")
    .csv("C:\inputs\TEST.txt") 
Run Code Online (Sandbox Code Playgroud)

但是,在这种情况下,Glue 不保证它们提供合适的 Spark 读取器。因此,如果您的错误与 CSV 缺少数据源有关,那么您应该通过--extra-jars参数提供指向其位置的 s3 路径,将spark-csv库添加到 Glue 作业中。


RK.*_*RK. 5

在以下 2 种情况下,我测试工作正常:

将文件从 S3 加载到 Glue。

dfnew = glueContext.create_dynamic_frame_from_options("s3", {'paths': ["s3://MyBucket/path/"] }, format="csv" )

dfnew.show(2)
Run Code Online (Sandbox Code Playgroud)

从已经通过 Glue Crawler 生成的 Glue 数据库和表加载数据。

DynFr = glueContext.create_dynamic_frame.from_catalog(database="test_db", table_name="test_table")
Run Code Online (Sandbox Code Playgroud)

DynFr 是一个 DynamicFrame,所以如果我们想在 Glue 中使用 Spark 代码,那么我们需要将它转换成一个普通的数据帧,如下所示。

df1 = DynFr.toDF()
Run Code Online (Sandbox Code Playgroud)