[
https://issues.apache.org/jira/browse/FLINK-21132?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=17274306#comment-17274306
]
Roman Khachatryan commented on FLINK-21132:
-------------------------------------------
I think it would be risky to change the behaviour of endOfInput. Existing
operators may rely on it being called even in case of stopWithSavepoint.
Instead, I'd propose to add a new interface StopWithSavepointAware \{ void
stoppedWithSavepoint(); } and call this method for operators implementing it
and endInput otherwise. Also in case of "true" end of input the old method
would be called.
Then IcebergFilesCommitter should implement it.
To distinguish true end-of-input from stop-with-savepoint we can rely purely on
checkpoint properties:
# Sources receive sync-savepoint RPC and only then call operatorChain.close /
endInput
# Downstreams must receive and confirm the checkpoint before getting EOP from
source
So in both cases it' enough to remember that a sync-savepoint was received.
Alternative approach is quite invasive as it in addition would change Task,
AbstractInvokable, passing and interprting EOP event.
I've prepared a prototype here:
[https://github.com/rkhachatryan/flink/tree/f21132]. It uses [~kezhuw] IT case
(thanks for providing it!).
It doesn't solve FLIP-27 issue but I think FLIP-27 should be addressed
separately.
WDYT?
> BoundedOneInput.endInput is called when taking synchronous savepoint
> --------------------------------------------------------------------
>
> Key: FLINK-21132
> URL: https://issues.apache.org/jira/browse/FLINK-21132
> Project: Flink
> Issue Type: Bug
> Components: Runtime / Task
> Affects Versions: 1.10.2, 1.10.3, 1.11.3, 1.12.1
> Reporter: Kezhu Wang
> Priority: Major
>
> [~elkhand](?) reported on project iceberg that {{BoundedOneInput.endInput}}
> was
> [called|https://github.com/apache/iceberg/issues/2033#issuecomment-765864038]
> when [stopping job with
> savepoint|https://github.com/apache/iceberg/issues/2033#issuecomment-765557995].
> I think it is a bug of Flink and was introduced in FLINK-14230. The
> [changes|https://github.com/apache/flink/pull/9854/files#diff-0c5fe245445b932fa83fdaf5c4802dbe671b73e8d76f39058c6eaaaffd9639faL577]
> rely on {{StreamTask.afterInvoke}} and {{OperatorChain.closeOperators}} will
> only be invoked after *end of input*. But that is not true long before after
> [FLIP-34: Terminate/Suspend Job with
> Savepoint|https://cwiki.apache.org/confluence/pages/viewpage.action?pageId=103090212].
> Task could enter state called
> [*finished*|https://github.com/apache/flink/blob/3a8e06cd16480eacbbf0c10f36b8c79a6f741814/flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/tasks/StreamTask.java#L467]
> after synchronous savepoint, that is an expected job suspension and stopping.
> [~sunhaibotb] [~pnowojski] [~roman_khachatryan] Could you help confirm this ?
> For full context, see
> [apache/iceberg#2033|https://github.com/apache/iceberg/issues/2033]. I have
> pushed branch
> [synchronous-savepoint-conflict-with-bounded-end-input-case|https://github.com/kezhuw/flink/commits/synchronous-savepoint-conflict-with-bounded-end-input-case]
> in my repository. Test case
> {{SavepointITCase.testStopSavepointWithBoundedInput}} failed due to
> {{BoundedOneInput.endInput}} called.
> I am also aware of [FLIP-147: Support Checkpoints After Tasks
> Finished|https://cwiki.apache.org/confluence/x/mw-ZCQ], maybe the three
> should align on what *finished* means exactly. [~kkl0u] [~chesnay]
> [~gaoyunhaii]
--
This message was sent by Atlassian Jira
(v8.3.4#803005)