【发布时间】:2020-12-30 11:11:38
【问题描述】:
我特别希望通过更新数据并将数据插入到包含大约 4 万亿条记录的 DeltaLake 基表中来优化性能。
环境: 火花 3.0.0 三角湖 0.7.0
在上下文中,这是关于通过 DeltaLake 制作增量表,我将按步骤对此进行更详细的总结:
- 创建基表(增量)
- 获取周期性数据
- 将数据添加到基表中
步骤 1 和 2 已经完成,但是添加数据的时候性能是出了名的慢,例如添加一个 9GB 的 csv 大约需要 6 个小时,这主要是因为 delta 需要为每次更新重写数据,它还需要" 读取 "数据库中的所有数据。
这个表也被分区(PARTITIONED BY)并存储在集群的GDFS(HDFS)中,以确保spark节点可以执行操作。
基表的字段(由#分配的基数):
- ID:标识符,# 10000
- 类型:字符串,# ~ 30
- LOCAL_DATE:记录的本地日期
- DATE_UTC:UTC 注册日期
- 值:注册表值
- YEAR:列计算 int #4
- MONTH:计算列 int #12
- DAY:计算列 int #31
由于一般搜索是按时间进行的,因此决定按 YEAR、MONTH、DAY 中的 LOCAL_DATE 列进行分区,由于其高基数,排除了按 ID 和 LOCAL_DATE 列进行分区,(出于性能目的更糟),最后加了TYPE,如下:
spark.sql(f"""
CREATE OR REPLACE TABLE {TABLE_NAME} (
ID INT,
FECHA_LOCAL TIMESTAMP,
FECHA_UTC TIMESTAMP,
TIPO STRING,
VALUE DOUBLE,
YEAR INT,
MONTH INT,
DAY INT )
USING DELTA
PARTITIONED BY (YEAR , MONTH , DAY, TIPO)
LOCATION '{location}'
""")
从现在开始,通过每 5 天定期添加这些大约 9 Gb 的 csv 文件来提供增量。目前 MERGE 操作如下:
spark.sql(f"""
MERGE INTO {BASE_TABLE_NAME}
USING {INCREMENTAL_TABLE_NAME} ON
--partitioned cols
{BASE_TABLE_NAME}.YEAR = {INCREMENTAL_TABLE_NAME}.YEAR AND
{BASE_TABLE_NAME}.MONTH = {INCREMENTAL_TABLE_NAME}.MONTH AND
{BASE_TABLE_NAME}.DAY = {INCREMENTAL_TABLE_NAME}.DAY AND
{BASE_TABLE_NAME}.TIPO = {INCREMENTAL_TABLE_NAME}.TIPO AND
{BASE_TABLE_NAME}.FECHA_LOCAL= {INCREMENTAL_TABLE_NAME}.FECHA_LOCALAND
{BASE_TABLE_NAME}.ID= {INCREMENTAL_TABLE_NAME}.ID
WHEN MATCHED THEN
UPDATE SET {BASE_TABLE_NAME}.VALUE= {INCREMENTAL_TABLE_NAME}.VALUE,
{BASE_TABLE_NAME}.TIPO= {INCREMENTAL_TABLE_NAME}.TIPO
WHEN NOT MATCHED THEN
INSERT *
""")
需要考虑的一些事实:
- 本次 MERGE 操作的时间为 6 小时
- 基表是根据 230GB csv 数据创建的(现在 55GB 在增量中!)
- spark 应用程序配置处于集群模式,具有以下参数
- 基础设施由 3 个节点、32 个内核和 250GB RAM 组成,但与其他现有应用程序相比,它占用的安全性要少大约 -50% 的资源。
Spark 应用
mode = 'spark: // spark-master: 7077'
# mode = 'local [*]'
spark = (SparkSession.builder.master (mode)
.appName ("SparkApp")
.config ('spark.cores.max', '45')
.config ('spark.executor.cores', '5')
.config ('spark.executor.memory', '11g')
.config ('spark.driver.memory', '120g')
.config ("spark.sql.shuffle.partitions", f "200") # 200 only for
200GB delta table reads
.config ("spark.storage.memoryFraction", f "0.8")
# DeltaLake configs
.config ("spark.jars.packages", "io.delta:delta-core_2.12:0.7.0")
.config ("spark.sql.extensions",
"io.delta.sql.DeltaSparkSessionExtension")
.config ("spark.sql.catalog.spark_catalog",
"org.apache.spark.sql.delta.catalog.DeltaCatalog")
# Delta optimization
.config ("spark.databricks.delta.optimizeWrite.enabled", "true")
.config ("spark.databricks.delta.retentionDurationCheck.enabled",
"false")
.getOrCreate ()
)
【问题讨论】:
标签: python apache-spark