Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

[FLINK-37475] Drop ChangelogNormalize for piping from upsert source to sink #26306

Merged
merged 1 commit into from
Mar 21, 2025

Conversation

dawidwys
Copy link
Contributor

What is the purpose of the change

Remove ChangelogNormalize if possible for INSERT INTO SELECT * FROM

Verifying this change

Added plan and semantic tests

Does this pull request potentially affect one of the following parts:

  • Dependencies (does it add or upgrade a dependency): (yes / no)
  • The public API, i.e., is any changed class annotated with @Public(Evolving): (yes / no)
  • The serializers: (yes / no / don't know)
  • The runtime per-record code paths (performance sensitive): (yes / no / don't know)
  • Anything that affects deployment or recovery: JobManager (and its components), Checkpointing, Kubernetes/Yarn, ZooKeeper: (yes / no / don't know)
  • The S3 file system connector: (yes / no / don't know)

Documentation

  • Does this pull request introduce a new feature? (yes / no)
  • If yes, how is the feature documented? (not applicable / docs / JavaDocs / not documented)

@dawidwys dawidwys marked this pull request as draft March 17, 2025 08:11
@flinkbot
Copy link
Collaborator

flinkbot commented Mar 17, 2025

CI report:

Bot commands The @flinkbot bot supports the following commands:
  • @flinkbot run azure re-run the last Azure build

@dawidwys dawidwys force-pushed the flink37475 branch 4 times, most recently from 2eb62e6 to eb97525 Compare March 17, 2025 18:53
@dawidwys dawidwys marked this pull request as ready for review March 18, 2025 07:26
Copy link
Contributor

@snuyanzin snuyanzin left a comment

Choose a reason for hiding this comment

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

I have only a couple of minor comments

@@ -577,6 +577,7 @@ void testFromAndToChangelogStreamUpsert() throws Exception {
tableEnv.fromChangelogStream(
changelogStream,
Schema.newBuilder().primaryKey("f0").build(),
// produce partial deletes
Copy link
Contributor

Choose a reason for hiding this comment

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

isn't the comment incorrect? it produces full deletes?

Copy link
Contributor Author

Choose a reason for hiding this comment

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

no, it's correct. I changed it that, by default upsert produces partial deletes.

@dawidwys
Copy link
Contributor Author

@flinkbot run azure

…o sink

* Add information about DELETE to ChangelogMode
* Adapt FlinkChangelogModeInferenceProgram to remove ChangelogNormalize
  if possible
@dawidwys dawidwys merged commit 50d46e1 into apache:master Mar 21, 2025
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.

None yet

5 participants