Skip to content

Spark: Added Scala/Java IcebergMergeInto Api on spark3.4 - #7607

Closed
giucris wants to merge 5 commits into
apache:mainfrom
giucris:feature/added-merge-into-api
Closed

giucris wants to merge 5 commits into
apache:mainfrom
giucris:feature/added-merge-into-api

Conversation

@giucris

@giucris giucris commented May 14, 2023

Copy link
Copy Markdown

Hi everyone.
Here a proposal to add on Iceberg an IcebergMergeInto dsl.

closes #3665

What has been added
A scala/java MergeInto Iceberg API.

What has been removed?
Nothing

What has benn changed?

  • A test utility method to easily create Dataset, previously it was a private ones

What is missing

  • Add to other spark versions
  • Complete other test cases
  • Docs enhancement

How it was developed

  • A new scala class has been added inside spark extensions. A simple builder that generates and execute an UnresolvedMergeIntoIcebergTable. The same one that is intercepted and rewritten during the execution of standard sql queries.

How it was tested

  • The same TestMerge tests were duplicated for test merge api

@github-actions github-actions Bot added the spark label May 14, 2023
@giucris giucris changed the title Spark: Added Scala/Java IcebergMergeInto Api on spark3.3 Spark: Added Scala/Java IcebergMergeInto Api on spark3.4 May 14, 2023
@giucris
giucris force-pushed the feature/added-merge-into-api branch from 1637c18 to 2ce216f Compare May 17, 2023 11:59
@giucris
giucris marked this pull request as ready for review May 17, 2023 12:00
@RussellSpitzer

RussellSpitzer commented May 17, 2023 •

Copy link
Copy Markdown
Member

First high level comment, I think we really need an underlying Java API as well. I think there should be a first class static class in java that can be used similar to the extension here. Basically I don't want to have the word "apply" in that invocation if we can help it. Other than that I think the API looks good at first glance, I assume it is matching the Delta api?

@giucris

giucris commented May 17, 2023

Copy link
Copy Markdown
Author

Hi @RussellSpitzer

The API is inspired by the delta ones to avoid confusion in those switching from one table format to another.

I think a Java Class with a static method that under calls the same builder might be enough to solve the "apply" problem and make everything look more "java" oriented.

I will try to add it

@Neuw84

Neuw84 commented Jun 8, 2023

Copy link
Copy Markdown
Contributor

Hi,

It would be great if we can push this feature to merge.

I could help if needed too on the Java side.

Let me know!

@RussellSpitzer

Copy link
Copy Markdown
Member

Other reviewers would be welcome, I'm a bit pressed for time

@giucris

giucris commented Jun 8, 2023

Copy link
Copy Markdown
Author

Sorry.
I was a bit underpressure on last 3 weeks.
I will try to complete as soon I have some time

@Neuw84

Neuw84 commented Jun 8, 2023

Copy link
Copy Markdown
Contributor

@giucris let me know if I can help on anything 👍

* @param source : Dataset[Row]
* @return [[IcebergMergeIntoBuilder]]
*/
def using(source: Dataset[Row]): IcebergMergeIntoBuilder =

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why is the target table using the table name while the source table using a Dataset? Should we support the table name as well?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ConeyLiu In my opinion the main use of a Merge api is to handle the merge between a source dataset and a target table. To avoid having to register the dataset as a temporary table and perform the operation in sql.

In case both the "delta" and the target table are already registered I don't see why use this api in favour of a simple SQL query.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hi @ConeyLiu can we consider this thread as solved?

*
*/
class IcebergMergeIntoBuilder(
private val targetTable: String,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

code style should be fixed

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hi @ConeyLiu it should be fixed now

* @param condition : [[Column]] expression
* @return [[IcebergMergeIntoBuilder]]
*/
def when(condition: Column): IcebergMergeIntoBuilder =

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The when seems a little confused with whenMatched / whenNotMatched. What about on?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ConeyLiu I agree.
I use when to have similarity with Delta MergeApi. But I completely agree that on could be better then when.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I would keep this as "on" because fundamentally this is setting shuffle join conditions (join on) while the other functions are only exercised at execution time.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hi @ConeyLiu and @RussellSpitzer this should be fixed now

this.whenMatchedActions :+ mergeAction,
this.whenNotMatchedActions
)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

blank line

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ConeyLiu this should be fixed now

*
*/
class IcebergMergeWhenMatchedBuilder(
private val icebergMergeIntoBuilder: IcebergMergeIntoBuilder,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

code style

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hi @ConeyLiu it should be fixed now

*
* It will update the target records accordingly with the assignment map provided as input.
*
* @param set : Map[String,String]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Maybe add doc for it

*
* It will update the target records accordingly with the assignment map provided as input.
*
* @param set : java.util.Map[String,String]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

same here

*
* It will update the target records accordingly with the assignment map provided as input.
*
* @param set : java.util.Map[String,String]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

same as here

*
* It will update the target records accordingly with the assignment map provided as input.
*
* @param set : java.util.Map[String,String]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

same as here

)
}

private def updateAction(set:Map[String,Column]): IcebergMergeIntoBuilder = {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

code tyle

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ConeyLiu this should be fixed now

condition,
set.map(x => Assignment(expr(x._1).expr, x._2.expr)).toSeq))
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

blank line

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ConeyLiu this should be fixed now

* Builder to specify IcebergMergeWhenNotMatched actions
*/
class IcebergMergeWhenNotMatchedBuilder(
private val icebergMergeIntoBuilder: IcebergMergeIntoBuilder,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

codestyle

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hi @ConeyLiu this should be fixed now

new IcebergMergeIntoBuilder(table, None, None, Seq.empty[MergeAction], Seq.empty[MergeAction])

/**
* Builder to specify an IcebergMergeInto action.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should move this to the above object?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ConeyLiu this should be fixed now

import scala.collection.JavaConverters


object IcebergMergeInto {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is there possible to provide the option functions to set the configs for the merge actions?

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This would be great, although I think we can't do that without a change of the actual Merge right?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hi @ConeyLiu I think it should be possible. I will try to configure it

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I did some research and I need more details on how you would like the merge api options to work.

  1. Do we want to set the options at the individual action level or at the global merge into level or both?

  2. What kind of options do we want to support? Runtime write config or Table properties or both?

To date and following the iceberg docs, the options really useful (IMHO) for merge, e.g., write.merge.mode, write.delete.mode, etc., are table properties and are currently set on the target table side during create or after altering the table.

However, there are some runtime configurations that could be provided at write time, these

  1. In case we also want to support table properties, what is expected after the merge operation with configured table properties? Target table containing the new properties or we want that the properties are used only in that runtime (In this case I don't know how feasible it is, I need to investigate further)

Given all these open points, can we addressing this need as a new feature, thus in a subsequent PR of the IcebergMergeInto API?

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@giucris did you conclude on this? I would like to pass write options in merge query which seems to be missing now.

}

private Dataset<Row> toDS(String schema, String jsonData) {
protected Dataset<Row> createDataset(String schema, String jsonData) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

why change the name

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ConeyLiu Previously this method was private. and the name toDS could fit. Now it has been aligned with the other utility methods that are protected and all this method begin with 'create' prefix

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I would agree that we shouldn't change method names if we can help it in a pr that is already quite large. I would just keep it as createDataset and do the change to toDS in another pr

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ConeyLiu @RussellSpitzer this should be fixed now

@ConeyLiu

ConeyLiu commented Jun 9, 2023

Copy link
Copy Markdown
Contributor

Thanks @giucris for the great work. left some API and code style comments.

@giucris

giucris commented Jun 9, 2023

Copy link
Copy Markdown
Author

@giucris let me know if I can help on anything 👍

@Neuw84 Feel free to contribute as you can

public synchronized void testMergeWithConcurrentTableRefresh() throws Exception {
// this test can only be run with Hive tables as it requires a reliable lock
// also, the table cache must be enabled so that the same table instance can be reused
Assume.assumeTrue(catalogName.equalsIgnoreCase("testhive"));

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not sure I understand why this can't be run with hadoop catalog. Should still be fine. The filesystem here should be the local fs and should be ok?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hi @RussellSpitzer. All the tests added are a "brutal" copy and paste of the tests inside the "sql" merge.

}

@Test
public synchronized void testMergeWithConcurrentTableRefresh() throws Exception {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not sure what the goal of this test is? We aren't really testing the API here but the actual underlying implementation of the merge. The test could fit in the merge implementation tests, I also would imagine that we would want to use Tasks.foreach or something like that to do our concurrent tasks, should be a bit clearer than the current semaphore and multi threaded approach?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hi @RussellSpitzer this test is a copy of the existing test inside the SQLMerge test cases.

I have started the development with all battery of test that already exist in the SQLMerge.

If this test does not make any sense I can drop it.
But this is already in the standard sql merge.

executorService.submit(
() -> {
for (int numOperations = 0; numOperations < Integer.MAX_VALUE; numOperations++) {
while (barrier.get() < numOperations * 2) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

integer overflow here

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hi @RussellSpitzer this test is a copy of the existing test inside the SQLMerge test cases.

I could drop it or just fix the test, but I would have to first get a good understanding of what is being tested, but it should also be fixed on the SQLMerge test cases side, perhaps in a later MR.

@giucris

giucris commented Sep 23, 2023

Copy link
Copy Markdown
Author

Hi @RussellSpitzer and @ConeyLiu.

First of all, thank you for your time and for your review and excuse me for coming back after a while. I've been a little busy.

I have been trying to fix all your threads. On others I opened a discussion because I need more clarity on how to proceed.

Now there is no need to have a Java API since I have included a table method as the main method to start the builder to avoid duplicate code.

IcebergMergeInto
   .table("icebergTable")
  .using(ds.as("source"))
  .on("source.id == icebergTable.id")
  .whenMatched()
  .updateAll()
  .whenNotMatched()
  .insertAll()
  .merge()

Let me know what you think, which threads can be closed, which ones to keep ongoing and which ones address into another PR.

Thank you again :)

@github-actions

Copy link
Copy Markdown

This pull request has been marked as stale due to 30 days of inactivity. It will be closed in 1 week if no further activity occurs. If you think that’s incorrect or this pull request requires a review, please simply write any comment. If closed, you can revive the PR at any time and @mention a reviewer or discuss it on the dev@iceberg.apache.org list. Thank you for your contributions.

@github-actions github-actions Bot added the stale label Aug 30, 2024
@github-actions

github-actions Bot commented Sep 7, 2024

Copy link
Copy Markdown

This pull request has been closed due to lack of activity. This is not a judgement on the merit of the PR in any way. It is just a way of keeping the PR queue manageable. If you think that is incorrect, or the pull request requires review, you can revive the PR at any time.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Adding MergeInto into the Spark Scala API

5 participants