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-34371][runtime] Support EndOfStreamTrigger and isOutputOnlyAfterEndOfStream operator attribute to optimize task deployment #24272
Conversation
...time/src/main/java/org/apache/flink/runtime/checkpoint/CheckpointCoordinatorDeActivator.java
Outdated
Show resolved
Hide resolved
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
@yunfengzhou-hub Thanks for the PR! Just left some comments below.
...-java/src/test/java/org/apache/flink/streaming/api/graph/StreamingJobGraphGeneratorTest.java
Outdated
Show resolved
Hide resolved
flink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/DefaultExecutionGraph.java
Outdated
Show resolved
Hide resolved
flink-runtime/src/main/java/org/apache/flink/runtime/scheduler/SchedulerBase.java
Outdated
Show resolved
Hide resolved
...-java/src/test/java/org/apache/flink/streaming/api/graph/StreamingJobGraphGeneratorTest.java
Outdated
Show resolved
Hide resolved
...ing-java/src/main/java/org/apache/flink/streaming/api/windowing/assigners/GlobalWindows.java
Outdated
Show resolved
Hide resolved
47bf841
to
479f495
Compare
Thanks for the comments @Sxnan @mohitjain2504 . I have updated the PR according to the comments. |
@yunfengzhou-hub Thanks for the update. LGTM! |
...g-java/src/main/java/org/apache/flink/streaming/api/operators/OperatorAttributesBuilder.java
Outdated
Show resolved
Hide resolved
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I've left some comments. PTAL. 😄
...-java/src/test/java/org/apache/flink/streaming/api/graph/StreamingJobGraphGeneratorTest.java
Show resolved
Hide resolved
...in/java/org/apache/flink/streaming/runtime/translators/OneInputTransformationTranslator.java
Show resolved
Hide resolved
9047c6a
to
b381679
Compare
...g-java/src/main/java/org/apache/flink/streaming/api/operators/OperatorAttributesBuilder.java
Outdated
Show resolved
Hide resolved
flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamNode.java
Show resolved
Hide resolved
...time/src/main/java/org/apache/flink/runtime/checkpoint/CheckpointCoordinatorDeActivator.java
Outdated
Show resolved
Hide resolved
flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/CheckpointFailureReason.java
Outdated
Show resolved
Hide resolved
...ava/src/main/java/org/apache/flink/streaming/api/transformations/PhysicalTransformation.java
Outdated
Show resolved
Hide resolved
b381679
to
f5c6bd8
Compare
…h outputOnlyAfterEndOfStream
39c5ae2
to
2bb2cc2
Compare
2bb2cc2
to
6567058
Compare
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
LGTM
What is the reason for closing? Did I miss something? |
@mohitjain2504 |
…or and StreamSortOperator This closes apache#24272
What is the purpose of the change
This pull request adds support for EndOfStreamTrigger and isOutputOnlyAfterEndOfStream operator attribute to optimize task deployment.
Brief change log
Verifying this change
Does this pull request potentially affect one of the following parts:
@Public(Evolving)
: yesDocumentation