Repository navigation
Feature Request: Support mergeSchema option when using Spark MERGE INTO #5556
Description
Activity
For reference, the issue here is that the user wants to be able to use
mergeSchermaoption when writing viaMERGE INTO.I'm not sure of a way to support that presently. If somebody does know, please comment 🙂
It would be great help if someone could help in achieving this functionality.. We are struggling to do this thing manually...
Just an FYI but I would update the title to be
Feature Request: Support mergeSchema option when using Spark MERGE INTO. This is more explicit and gets to the heart of what it is you need.The hints might not be something we can add without changing Spark, but the core of the idea is that you need
mergeSchemato work with MERGE INTO (which is currently SQL only).Removing the implementation constraint from the title might attract more eyeballs / bring more ideas to the table (as ultimately you don’t care about anything other than needing
mergeIntoto work).Also, what about a table property? Does the table experience writes where you explicitly do not want
mergeSchema?Generally, I think
mergeSchemis safer as a per-query option and is somewhat unsafe as a table level configuration. But Spark makes it somehat hard to support that as there’s no Dataframe support forMERGE INTOcurrently.For the long run, I’m going to bring up adding a merge into API to the dataframe / dataset API in Spark. But that could take a while. We might be able to provide implicit classes so that it’s do-able using the dataframe API in just Iceberg, but in the long run that should be moved to Spark (though that doesn’t solve your immediate problem, I know).
- changed the title
[-]Support Hints for Dataframe Writer Options Like 'mergeSchema'[/-][+]Feature Request: Support mergeSchema option when using Spark MERGE INTO[/+]on Aug 18, 2022 This issue has been automatically marked as stale because it has been open for 180 days with no activity. It will be closed in next 14 days if no further activity occurs. To permanently prevent this issue from being considered stale, add the label 'not-stale', but commenting on the issue is preferred when possible.
This issue has been closed because it has not received any activity in the last 14 days since being marked as 'stale'
Not stale
Reacted by Diego Tobon, clamar14, baunz, Daniel, stu-par, jadamsoneko, yasuyuki okasaka, advaitpathak-cod, Nicolas Ferrario and Igor BermanThis is definitely something we'd be interested in - We are doing something similar with Glue and would like to be able to support schema evolution with a
MERGE INTOUPSERT. Currently we have to manually modify the iceberg schema every time our source schema changes.This is a very similar architecture to what we're doing - https://aws.amazon.com/blogs/big-data/automate-replication-of-relational-sources-into-a-transactional-data-lake-with-apache-iceberg-and-aws-glue/
Reacted by Andrea Campolonghi, K Faiqoh, Diego Tobon and jadamsonekoReopening due to interest in this
Reacted by Jared Bates, Andrea Campolonghi, kongul, Itagyba Abondanza Kuhlmann, Diego Tobon, Matt Corley, Daryna, stu-par, jadamsoneko and advaitpathak-cod@kbendick is this still on your radar? If not could you give me some direction on where I could start to look at.
Any updates on this feature? I also have a strong interest in iceberg providing this solution.
Anyone who would like to work on the issue is welcome to, there is currently no one I know working on it.
Delta Lake has the ability to set
spark.databricks.delta.schema.autoMerge.enabled. I find this approach interesting as it can be used only when required. Once set the automatic schema evolution works for every write operation.Reacted by Fokko Driesprong and Karthik PrabhakarIs it still being worked on? It would be nice if we can have either:
- Schema evolution (schema merge) for merge sql statements
or
2 DataFrame API for merge queries
- Schema evolution (schema merge) for merge sql statements
Man if they added this, Iceberg would be used even more. This is the biggest pain point of Iceberg and I hope this is sorted because I don't see hard it could be, compared to what is already done by Iceberg, to add this feature.
With the code for MERGE being moved to Spark with Spark 3.5 & Iceberg 1.5.0, is it even possible to support this feature solely with changes in Iceberg, or are changes in Spark required?
I am interested in digging more into this feature and working on it if no one has started.
Reacted by Lorenzo Coacci, Tartuger, Michael Dingess and [테크타카] 김용집I have done some research into this and adding my notes below. Please correct me if any part is incorrect.
I believe this issue will be fairly easy to support in Iceberg after Spark 4.0 is released and Iceberg is integrated with it. Spark 4.0 dataframe API supports MergeInto and the dataframe API also supports schema evolution.
Past work has attempted to add the MergeInto dataframe API into Iceberg to provide this functionality earlier than Spark 4.0 timeline but this is a sizable change and has not been completed.
Spark 4.0 is planned for release around June 2025 and Iceberg integration may come soon after that, so if my understanding is correct then this issue may be resolved in mid to late 2025.
Spark support for DSV2 MERGE INTO WITH SCHEMA EVOLUTION will be released in Spark 4.1, see apache/spark#51698 and follow ups
Reacted by Zach All and Nguyễn Quốc Vương
Feature Request / Improvement
Hi Team,
I am using Iceberg in my project and I found a big thing which is missing from Iceberg which is easily available in Apache Hudi and Deltalake that is "merge schema". If possible this feature need to added into the Iceberg. I am attaching my last ticket which is explaining the problem that I am facing.Please find the below ticket for the refrence.
#5548
@rdblue any thoughts on this?
Query engine
Spark