【发布时间】:2020-02-12 04:30:45
【问题描述】:
给定一个应用程序将 csv 转换为 parquet(从和到 S3),几乎没有转换:
for table in tables:
df_table = spark.read.format('csv') \
.option("header", "true") \
.option("escape", "\"") \
.load(path)
df_one_seven_thirty_days = df_table \
.filter(
(df_table['date'] == fn.to_date(fn.lit(one_day))) \
| (df_table['date'] == fn.to_date(fn.lit(seven_days))) \
| (df_table['date'] == fn.to_date(fn.lit(thirty_days)))
)
for i in df_one_seven_thirty_days.schema.names:
df_one_seven_thirty_days = df_one_seven_thirty_days.withColumnRenamed(i, colrename(i).lower())
df_one_seven_thirty_days.createOrReplaceTempView(table)
df_sql = spark.sql("SELECT * FROM "+table)
df_sql.write \
.mode("overwrite").format('parquet') \
.partitionBy("customer_id", "date") \
.option("path", path) \
.saveAsTable(adwords_table)
我在使用 spark EMR 时遇到了困难。
在本地使用 spark 提交,运行没有困难(140MB 数据)并且速度非常快。 但在 EMR 上,情况就另当别论了。
第一个“adwords_table”将毫无问题地转换,但第二个保持空闲状态。
我浏览了 EMR 提供的 spark 作业 UI,我注意到一旦完成此任务:
列出 187 个路径的叶子文件和目录:
20 分钟后没有任何反应。所有任务都处于“已完成”状态,没有新任务开始。 我正在等待 saveAsTable 启动。
我的本地机器是 8 核 15GB,集群由 10 个节点组成 r3.4xlarge: 32 vCore、122 GiB 内存、320 SSD GB 存储 EBS 存储:200 GiB
配置使用maximizeResourceAllocation true,我只将 --num-executors / --executor-cores 更改为 5
有人知道为什么集群进入“空闲”状态并且没有完成任务吗? (它最终会在 3 小时后无错误地崩溃)
编辑: 通过删除所有胶水目录连接 + 降级使用 hadoop,我取得了一些进展:hadoop-aws:2.7.3
现在 saveAsTable 工作正常,但是一旦完成,我看到执行程序被删除并且集群处于空闲状态,该步骤没有完成。
所以我的问题还是一样。
【问题讨论】:
-
您尝试写入的确切 s3 路径是什么?
-
path = "s3a://{0}/{1}/{2}".format(S3_DESTINATION_RAW_BUCKET, S3_PROCESSED_ADWORDS_PATH, adwords_table)
-
您是否尝试过单独运行它们而不是循环运行
标签: apache-spark hadoop amazon-emr