【问题标题】:How to improve the performance of a merge operation with an incremental DeltaLake table?如何使用增量 DeltaLake 表提高合并操作的性能?
【发布时间】:2020-12-30 11:11:38
【问题描述】:

我特别希望通过更新数据并将数据插入到包含大约 4 万亿条记录的 DeltaLake 基表中来优化性能。

环境: 火花 3.0.0 三角湖 0.7.0

在上下文中,这是关于通过 DeltaLake 制作增量表,我将按步骤对此进行更详细的总结:

  1. 创建基表(增量)
  2. 获取周期性数据
  3. 将数据添加到基表中

步骤 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


    【解决方案1】:

    我选择分享这个答案,以便您可以利用一些技巧。

    Delta 建议使用所有分区列,这样最终的数据搜索会更少,由“剪枝”的效果给出

    因此有必要确定合并可以更新数据的所有情况,为此 对增量数据进行查询以生成此类型的字典:

    filter_columns = spark.sql (f "" "
    SELECT
        YEAR,
        MONTH,
        DAY,
        COLLECT_LIST (DISTINCT TYPE) AS TYPES
    Incremental FROM
    GROUP BY YEAR, MONTH, DAY
    ORDER BY 1, 2, 3
    "" ") .toPandas ()
    

    有了这个df,就可以生成合并必须更新/插入的条件:

    [! [df按年、月、日、类型分组]1]1

    然后它生成了一个名为“final_cond”的字符串,如下所示:

    dic = filter_columns.groupby (['YEAR', 'MONTH', 'DAY']) ['TYPE']. apply (lambda grp: list (grp.value_counts (). index)). to_dict ()
    final_cond = ''
    index = 0
    for key, value in dic.items ():
        cond = ''
        year = key [0]
        month = key [1]
        day = key [2]
        variables = ','. join (["'" + str (x) + "'" for x in value [0]])
        or_cond = '' if index + 1 == len (dic) else '\ nOR \ n'
        
        cond = f "" "({BASE_TABLE_NAME} .YEAR == {year} AND {BASE_TABLE_NAME} .MONTH == {month} AND {BASE_TABLE_NAME} .DAY == {day} AND {BASE_TABLE_NAME}. TYPE IN ({variables} )) "" "
          
        final_cond = final_cond + cond + f '{or_cond}'
        index + = 1
        #break
        
    print (final_cond)
    

    [! [字符串条件]2]

    最后我们将这些条件添加到 MERGE 中:

    ...
    WHEN MATCHED AND ({final_cond}) THEN
    ...
    

    这个简单的“过滤器”减少了大型操作的合并时间

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-10-29
      • 2016-11-01
      • 2018-04-06
      • 2017-04-13
      • 1970-01-01
      • 2011-08-27
      • 1970-01-01
      • 2022-06-16
      相关资源
      最近更新 更多