【发布时间】:2021-09-01 21:51:35
【问题描述】:
我决定在这里发布这个问题,因为我没有进一步的想法。我一直在尝试以多种不同方式提高代码的性能,但仍然面临同样的问题。
从一开始...
-
我正在从 700mb 的 csv 文件(约 600 万行)创建数据帧 + 缓存这些数据以供进一步处理
-
我正在删除基于 5 列的重复项,这为我提供了新的更轻的数据框(约 300 万行)
-
我正在通过 jdbc spark 连接器(允许批量输入)将新的更轻量级数据帧输入到我的 azure sql db。
输入过程大约需要 9 分钟才能完成,而前一个过程需要不到一秒。我正在以多种不同的方式重新分区我的 DF,但仍然没有。
我正在使用来自 azure 的高级试用版数据块。它只允许我创建一种类型的集群 -> 单节点,14gb 内存,4 个核心。
我对 spark 比较陌生,在我的理解中 4 个内核意味着 4 个真正的分区?那么这个集群可能是它需要这么长时间的原因吗?或者也许我遗漏了一些东西,但仍然可以提高处理速度。
这是我的代码:
第一部分:
from pyspark.sql.types import StructType, StructField, StringType, DateType, DecimalType
weather_schema = StructType([
StructField("EventId", StringType(), True),
StructField("Type", StringType(), True),
StructField("Severity", StringType(), True),
StructField("StartTime(UTC)", DateType(), True),
StructField("EndTime(UTC)", DateType(), True),
StructField("TimeZone", StringType(), True),
StructField("AirportCode", StringType(), True),
StructField("LocationLat", DecimalType(18,0), True),
StructField("LocationLng", DecimalType(18,0), True),
StructField("City", StringType(), True),
StructField("Country", StringType(), True),
StructField("State", StringType(), True),
StructField("ZipCode", StringType(), True)
])
weatherDF = (spark.read
.option("sep", ",")
.option("header", True)
.schema(weather_schema)
.csv(source+"weather.csv"))
weatherDF.cache()
weatherDF.show()
命令耗时 0.47 秒
第二部分:
from pyspark.sql.functions import count, col, asc
clearedDF = (weatherDF.dropDuplicates(["Type", "Severity", "StartTime(UTC)", "EndTime(UTC)", "City"])
.orderBy(col("EventId").asc())
)
命令耗时 0.03 秒
第三部分(有问题):
jdbcHost = dbutils.secrets.get(scope="adminSecrets", key="dbServer")
jdbcPort = "1433"
jdbcDatabase = dbutils.secrets.get(scope="adminSecrets", key="dbName")
jdbcUrl = "jdbc:sqlserver://{0}:{1};database={2}".format(jdbcHost, jdbcPort, jdbcDatabase)
properties = {
"user":dbutils.secrets.get(scope="adminSecrets", key="adminLogin"),
"password":dbutils.secrets.get(scope="adminSecrets", key="adminPassword"),
"driver" : "com.microsoft.sqlserver.jdbc.spark"
}
try:
(clearedDF.repartition(4).write
.format(properties["driver"])
.option("batchsize", 100000)
.option("tableLock", "true")
.option("url", jdbcUrl)
.option("dbtable", "dbo.weather")
.option("user", properties["user"])
.option("password", properties["password"])
.mode("overwrite")
.save()
)
except ValueError as error:
print(str(error))
命令耗时 8.57 分钟
【问题讨论】:
-
你真的需要
repartition(4)吗?但是,是的,数十次访问外部服务(数百万行,除以批量大小)将需要一些时间 -
@OneCricketeer 如果我错了,请纠正我。在不添加 repartition(4) 的情况下,它需要相同的时间。意思是spark默认使用所有可用资源(核心)?
-
假设您使用的是
master("local[*]"),那么它将使用所有内核。不过,数据框本身可能有更多分区 -
所以如果我有 4 个核心并且数据帧有更多的分区,那么这意味着 4 个核心有 4 个分区,而这 4 个分区有自己的分区?我为混乱的解释道歉,但我相信你会明白我的意思。
标签: apache-spark pyspark apache-spark-sql databricks azure-databricks