【发布时间】:2018-09-10 10:38:53
【问题描述】:
我正在尝试从 Oracle DB 中提取数据并使用 Apache Spark 2.3.1 将其放入 AWS S3。这项工作一直运行良好,直到最后一个阶段并被卡在那里。我不认为数据有偏差,因为每个阶段都有相同数量的记录。下面是我在 spark 中使用的查询。
url = "jdbc:oracle:thin:@IP:PORT/SID"
user = "user"
password = "password"
driver = "oracle.jdbc.driver.OracleDriver"
table = "table"
fetchSize = 1000
partitionColumn = "num_rows"
date1 = (datetime.today() - td(days=42)).date().strftime('%d-%b-%Y')
date2 = (datetime.today() - td(days=2)).date().strftime('%d-%b-%Y')
query = "(select min(rownum) as min, max(rownum) as max from "+table+" where date>='"+str(date1)+"' and date<='"+str(date2)+"') tmp1"
print(query)
DF = spark.read.format("jdbc").option("url", url) \
.option("dbtable", query) \
.option("user", user) \
.option("password", password) \
.option("driver", driver) \
.load()
lower_bound, upper_bound = DF.first()
lower_bound = int(lower_bound)
upper_bound = int(upper_bound)
numPartitions = int(upper_bound/fetchSize)+1
print(lower_bound,upper_bound)
print(numPartitions)
query = "(select t1.*, ROWNUM as num_rows from (select * from " + table + " where date>='"+str(date1)+"' and date<='"+str(date2)+"') t1) tmp2"
print(query)
DF = spark.read.format("jdbc").option("url", url) \
.option("dbtable", query) \
.option("user", user) \
.option("password", password) \
.option("fetchSize",fetchSize) \
.option("numPartitions", numPartitions) \
.option("partitionColumn", partitionColumn) \
.option("lowerBound", lower_bound) \
.option("upperBound", upper_bound) \
.option("driver", driver) \
.load()
path = "s3://my_path"
DF.write.mode("overwrite").parquet(path)
代码基本上是提取最近 42 天的数据并将其放入 S3 存储桶中。下面是直到写入语句的输出。代码在 '10-Sep-2018'
上运行(select min(rownum) as min, max(rownum) as max from table where date>='30-Jul-2018' and date<='08-Sep-2018') tmp1
(1, 2195427)
2196
(select t1.*, ROWNUM as num_rows from (select * from table where date>='30-Jul-2018' and date<='08-Sep-2018') t1) tmp2
如你所见,
- 总记录数=2195427
- 每个分区的记录 = 1000
- 分区数 = 2196
所以该作业有 2196 个阶段,每个阶段提取 1000 条记录。这项工作在 2191/2196 卡住了,还有 5 个阶段要完成。
硬件规格:
我正在使用 r4.xlarge 机器。我的集群是 r4.xlarge 的 1 个 Master,2 个 Slaves。以下是我的驱动程序和执行程序规格。
spark.driver.cores 8
spark.driver.memory 24g
spark.driver.memoryOverhead 3072M
spark.executor.cores 1
spark.executor.memory 3g
spark.executor.memoryOverhead 512M
spark.yarn.am.cores 1
spark.yarn.am.memory 3g
spark.yarn.am.memoryOverhead 512M
第 1 至 2191 阶段在 1.3 小时内完成,但其余 5 个阶段卡住了三个多小时。
请在此处找到日志: https://github.com/rinazbelhaj/stackoverflow/blob/master/Spark_Log_10_Sept_2018
我无法找出这个问题的根本原因。
【问题讨论】:
-
嗨!我的一个朋友在这里报告了一个非常相似的问题 (stackoverflow.com/questions/54315191/…)。你找到解决办法了吗?
-
嗨!你找到解决方案了吗?我也面临同样的问题...
标签: apache-spark pyspark apache-spark-sql