【发布时间】:2019-08-16 02:46:35
【问题描述】:
有没有办法从 RedShift 的 tempDir 转储创建 DataFrame?
我的用例是当作业失败时,我想重试,但从转储到 S3 的临时数据转储继续,而不是从 RedShift 重新获取数据集,这是巨大的!
加载代码会这样做
val df1 = spark.read
.format("com.databricks.spark.redshift")
.option("url", jdbcUrl)
.option("dbtable", spmeTable)
.option("tempdir", tempDir)
.option("user", jdbcUsername)
.option("password", jdbcPassword)
.option("forward_spark_s3_credentials", true)
.load();
稍后的作业失败,但我想重新创建 df1 而不再次从 RedShift 获取任何内容。
有没有办法做到这一点?
在SparkSession下找到了一个名为createDataFrame的方法,不知道这是否可行……
https://spark.apache.org/docs/2.3.0/api/java/org/apache/spark/sql/SparkSession.html
更新 #1
临时目录看起来像这里的目录结构 https://docs.aws.amazon.com/redshift/latest/dg/r_UNLOAD_command_examples.html
我从 S3 打开了一个临时文件,它是用管道分隔的
edd66540-fa17-599b-9b22-7df29a5f9229|kNOCugU4wuKAUw7m2UXS7MfX|2018-11-27 19:48:44|POST|f|@NULL@|@NULL@|@NULL@|@NULL@|https://www.example.com/r/conversations/0grt6540-
更新 #2
据此 https://github.com/databricks/spark-redshift/tree/master/tutorial
将文件写入 S3 后,将使用自定义 InputFormat (com.databricks.spark.redshift.RedshiftInputFormat) 并行使用文件。该类类似于 Hadoop 的标准 TextInputFormat 类,其中键是文件中每一行开头的字节偏移量。然而,值类是 Array[String] 类型(与 TextInputFormat 不同,它的类型是 Text)。这些值是通过使用默认分隔符 (|) 拆分行来创建的。 RedshiftInputFormat 逐行处理 S3 文件以生成 RDD。然后将之前获得的模式应用于此 RDD,以将字符串转换为适当的数据类型并生成 DataFrame。
除了跳过卸载之外,知道怎么做吗?
【问题讨论】:
标签: apache-spark apache-spark-sql amazon-redshift databricks