【问题标题】:Spark, delta lake auto schema evolution for nested columnsSpark,Delta Lake 嵌套列的自动模式演变
【发布时间】:2021-05-25 20:14:53
【问题描述】:

合并时模式演化的深度是多少?

在以下情况下合并时,自动架构演变不起作用。

import json
d1 = {'a':'b','b':{'c':{'1':1}}}
d2 = {'a':'s','b':{'c':{'1':2,'2':2}}}
d3 = {'a':'v','b':{'c':{'1':4}}}

df1 = spark.read.json(spark.sparkContext.parallelize([json.dumps(d1)]))

#passes
df1.write.saveAsTable('test_table4',format='delta',mode='overwrite', path=f"hdfs://hdmaster:9000/dest/test_table4")


df2 = spark.read.json(spark.sparkContext.parallelize([json.dumps(d2)]))
df2.createOrReplaceTempView('updates')

query = """
MERGE INTO test_table4 existing_records 
        USING updates updates 
        ON existing_records.a=updates.a
        WHEN MATCHED THEN UPDATE SET * 
        WHEN NOT MATCHED THEN INSERT *
"""
spark.sql("set spark.databricks.delta.schema.autoMerge.enabled=true")
spark.sql(query) #passes



df3 = spark.read.json(spark.sparkContext.parallelize([json.dumps(d3)]))

df3.createOrReplaceTempView('updates')
query = """
MERGE INTO test_table4 existing_records 
        USING updates updates 
        ON existing_records.a=updates.a
        WHEN MATCHED THEN UPDATE SET * 
        WHEN NOT MATCHED THEN INSERT *
"""
spark.sql("set spark.databricks.delta.schema.autoMerge.enabled=true")
spark.sql(query) #FAILS #FAILS

当深度大于 2 并且传入的 df 缺少列时,这看起来会失败。 这是故意这样吗? 如果想追加,可以使用option("mergeSchema", "true") 完美处理。但我想 UPSERT 数据。但 Merge 无法处理此架构更改

使用 Delta Lake 版本 0.8.0

【问题讨论】:

    标签: apache-spark pyspark bigdata delta-lake


    【解决方案1】:

    在 Delta 0.8 中,这应该通过将spark.databricks.delta.schema.autoMerge.enabled 设置为true 来进行调节,除了mergeSchema 更适合append 模式。

    有关此功能的更多详细信息,请参阅Delta 0.8 announcement blog post

    【讨论】:

    • 嘿,我试过了。当您在嵌套 json 的更深层次上删除字段时,这不起作用
    • 但是当在共享代码中提到的根级别删除该字段时,删除工作有效
    猜你喜欢
    • 1970-01-01
    • 2019-10-06
    • 2021-07-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-10-15
    • 1970-01-01
    • 2020-12-02
    相关资源
    最近更新 更多