【发布时间】:2021-09-04 21:34:57
【问题描述】:
我有一个简单的 spark 工作,它做了 3 件事:
- 从 AWS S3 逐月读取 json 数据(数据按日期分区)。
- 对数据做一些最少的处理。
- 将处理后的月度数据覆盖到同一来源。
作业成功处理并覆盖了几个月,但在覆盖到 S3 时随机引发此异常:
原因:com.amazon.ws.emr.hadoop.fs.shaded.com.amazonaws.services.s3.model.AmazonS3Exception:找不到一个或多个指定部件。该部分可能尚未上传,或者指定的实体标签可能与该部分的实体标签不匹配。 (服务:Amazon S3;状态代码:400;错误代码:InvalidPart;请求 ID:123;S3 扩展请求 ID:xyz/ck+foo/bar=)
这里是 PySpark Job 的代码 sn-p:
spark = SparkSession.builder.appName('simple_app').getOrCreate()
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
source_data_lake_path = "s3://my-data-lake/data"
months_not_found, months_cleaned = [], []
for month in ['2020-01-*', '2020-02-*', '2020-03-*', ....]:
try:
data = spark.read.json(f"{source_data_lake_path}/persist_date={month}")
except AnalysisException:
months_not_found.append(month)
continue
cleaned_data, cleaned = clean_data(data)
if cleaned:
cleaned_data = cleaned_data.withColumn("persist_date", F.to_date(F.col("persist_timestamp")))
cleaned_data.repartition("persist_date").write.partitionBy("persist_date").mode("overwrite").json(
source_data_lake_path
)
months_cleaned.append(month)
我的发现:
- 使用 s3://,因为 s3:// 和 s3n:// 在 AWS EMR 的上下文中可以互换,而 s3a:// 与 EMR 不兼容。
- 我认为这是因为多个并发写入并减少了节点数。由于此异常,作业有时仍会随机失败。
【问题讨论】:
-
多少个分区?您的分段上传似乎失败了。 S3 有 api 限制,例如每秒 3500 个相同前缀 PUT 请求的请求。
-
@Lamanus 我不知道这是否是您问题的正确答案,但每个日期都有一个分区,所以大约。每月30。有什么方法可以保持在 API 限制范围内并且仍然能够覆盖所有这些数据?
标签: amazon-web-services apache-spark amazon-s3 pyspark