[
https://issues.apache.org/jira/browse/FLINK-18398?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=17154322#comment-17154322
]
Jun Qin commented on FLINK-18398:
---------------------------------
To summarize what happened:
At 2020-06-15 19:52:03.664Z, operator {{graph18}} failed to complete snapshot
1067096 due to an error in ElasticsearchSink:
{code:java}
# elastic_tm_log.txt
2020-06-15 19:52:03.664Z ERROR [ I/O dispatcher 229]
.f.s.c.e.ElasticsearchSinkBase : Failed Elasticsearch bulk request: request
retries exceeded max retry timeout [30000]java.io.IOException: request retries
exceeded max retry timeout [30000]2020-06-15 19:52:03.664Z ERROR [ I/O
dispatcher 229] .f.s.c.e.ElasticsearchSinkBase : Failed Elasticsearch bulk
request: request retries exceeded max retry timeout [30000]java.io.IOException:
request retries exceeded max retry timeout [30000] at
org.elasticsearch.client.RestClient$1.retryIfPossible(RestClient.java:411) at
org.elasticsearch.client.RestClient$1.failed(RestClient.java:398) at
org.apache.http.concurrent.BasicFuture.failed(BasicFuture.java:134) at
org.apache.http.impl.nio.client.AbstractClientExchangeHandler.failed(AbstractClientExchangeHandler.java:419)
at
org.apache.http.nio.protocol.HttpAsyncRequestExecutor.timeout(HttpAsyncRequestExecutor.java:375)
at
org.apache.http.impl.nio.client.InternalIODispatch.onTimeout(InternalIODispatch.java:92)
at
org.apache.http.impl.nio.client.InternalIODispatch.onTimeout(InternalIODispatch.java:39)
at
org.apache.http.impl.nio.reactor.AbstractIODispatch.timeout(AbstractIODispatch.java:175)
at
org.apache.http.impl.nio.reactor.BaseIOReactor.sessionTimedOut(BaseIOReactor.java:263)
at
org.apache.http.impl.nio.reactor.AbstractIOReactor.timeoutCheck(AbstractIOReactor.java:492)
at
org.apache.http.impl.nio.reactor.BaseIOReactor.validate(BaseIOReactor.java:213)
at
org.apache.http.impl.nio.reactor.AbstractIOReactor.execute(AbstractIOReactor.java:280)
at
org.apache.http.impl.nio.reactor.BaseIOReactor.execute(BaseIOReactor.java:104)
at
org.apache.http.impl.nio.reactor.AbstractMultiworkerIOReactor$Worker.run(AbstractMultiworkerIOReactor.java:588)
at java.lang.Thread.run(Thread.java:748)
2020-06-15 19:52:03.665Z ERROR [ I/O dispatcher 229] .e.g.i.s.Handler : index
request failure, index=index1, id=event1_2760347e-af6a-41b4-bd10-cdbfe707b547,
errorMessage=request retries exceeded max retry timeout [30000]2020-06-15
19:52:03.704Z INFO [...Sink (1/1)] f.s.a.o.AbstractStreamOperator : Could not
complete snapshot 1067096 for operator graph18
(1/1).java.lang.RuntimeException: An error occurred in ElasticsearchSink. at
org.apache.flink.streaming.connectors.elasticsearch.ElasticsearchSinkBase.checkErrorAndRethrow(ElasticsearchSinkBase.java:381)
at
org.apache.flink.streaming.connectors.elasticsearch.ElasticsearchSinkBase.checkAsyncErrorsAndRequests(ElasticsearchSinkBase.java:386)
at
org.apache.flink.streaming.connectors.elasticsearch.ElasticsearchSinkBase.snapshotState(ElasticsearchSinkBase.java:318)
at
org.apache.flink.streaming.util.functions.StreamingFunctionUtils.trySnapshotFunctionState(StreamingFunctionUtils.java:118)
at
org.apache.flink.streaming.util.functions.StreamingFunctionUtils.snapshotFunctionState(StreamingFunctionUtils.java:99)
at
org.apache.flink.streaming.api.operators.AbstractUdfStreamOperator.snapshotState(AbstractUdfStreamOperator.java:90)
at
org.apache.flink.streaming.api.operators.AbstractStreamOperator.snapshotState(AbstractStreamOperator.java:402)
at
org.apache.flink.streaming.runtime.tasks.StreamTask$CheckpointingOperation.checkpointStreamOperator(StreamTask.java:1420)
at
org.apache.flink.streaming.runtime.tasks.StreamTask$CheckpointingOperation.executeCheckpointing(StreamTask.java:1354)
at
org.apache.flink.streaming.runtime.tasks.StreamTask.checkpointState(StreamTask.java:991)
at
org.apache.flink.streaming.runtime.tasks.StreamTask.lambda$performCheckpoint$5(StreamTask.java:887)
at
org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$SynchronizedStreamTaskActionExecutor.runThrowing(StreamTaskActionExecutor.java:94)
at
org.apache.flink.streaming.runtime.tasks.StreamTask.performCheckpoint(StreamTask.java:860)
at
org.apache.flink.streaming.runtime.tasks.StreamTask.triggerCheckpoint(StreamTask.java:793)
at
org.apache.flink.streaming.runtime.tasks.StreamTask.lambda$triggerCheckpointAsync$3(StreamTask.java:777)
at java.util.concurrent.FutureTask.run(FutureTask.java:266) at
org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$SynchronizedStreamTaskActionExecutor.run(StreamTaskActionExecutor.java:87)
at org.apache.flink.streaming.runtime.tasks.mailbox.Mail.run(Mail.java:78) at
org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.processMail(MailboxProcessor.java:261)
at
org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:186)
at
org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:487)
at
org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:470)
at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:707) at
org.apache.flink.runtime.taskmanager.Task.run(Task.java:532) at
java.lang.Thread.run(Thread.java:748)Caused by: com.xxx.xxx.Exception: request
retries exceeded max retry timeout [30000] at
com.xxx.xxx.ExceptionFactory.create(ExceptionFactory.java:22) at
com.xxx.xxx.Handler.onFailure(Handler.java:18) at
org.apache.flink.streaming.connectors.elasticsearch.ElasticsearchSinkBase$BulkProcessorListener.afterBulk(ElasticsearchSinkBase.java:434)
at
org.elasticsearch.action.bulk.BulkRequestHandler$1.onFailure(BulkRequestHandler.java:77)
at org.elasticsearch.action.bulk.Retry$RetryHandler.onFailure(Retry.java:128)
at
org.elasticsearch.client.RestHighLevelClient$1.onFailure(RestHighLevelClient.java:606)
at
org.elasticsearch.client.RestClient$FailureTrackingResponseListener.onDefinitiveFailure(RestClient.java:629)
at org.elasticsearch.client.RestClient$1.retryIfPossible(RestClient.java:412)
at org.elasticsearch.client.RestClient$1.failed(RestClient.java:398) at
org.apache.http.concurrent.BasicFuture.failed(BasicFuture.java:134) at
org.apache.http.impl.nio.client.AbstractClientExchangeHandler.failed(AbstractClientExchangeHandler.java:419)
at
org.apache.http.nio.protocol.HttpAsyncRequestExecutor.timeout(HttpAsyncRequestExecutor.java:375)
at
org.apache.http.impl.nio.client.InternalIODispatch.onTimeout(InternalIODispatch.java:92)
at
org.apache.http.impl.nio.client.InternalIODispatch.onTimeout(InternalIODispatch.java:39)
at
org.apache.http.impl.nio.reactor.AbstractIODispatch.timeout(AbstractIODispatch.java:175)
at
org.apache.http.impl.nio.reactor.BaseIOReactor.sessionTimedOut(BaseIOReactor.java:263)
at
org.apache.http.impl.nio.reactor.AbstractIOReactor.timeoutCheck(AbstractIOReactor.java:492)
at
org.apache.http.impl.nio.reactor.BaseIOReactor.validate(BaseIOReactor.java:213)
at
org.apache.http.impl.nio.reactor.AbstractIOReactor.execute(AbstractIOReactor.java:280)
at
org.apache.http.impl.nio.reactor.BaseIOReactor.execute(BaseIOReactor.java:104)
at
org.apache.http.impl.nio.reactor.AbstractMultiworkerIOReactor$Worker.run(AbstractMultiworkerIOReactor.java:588)
... 1 common frames omittedCaused by: java.io.IOException: request retries
exceeded max retry timeout [30000] at
org.elasticsearch.client.RestClient$1.retryIfPossible(RestClient.java:411) ...
14 common frames omitted
{code}
Then the job cancellation/recovering was triggered due to Exceeded checkpoint
tolerable failure threshold:
{code:java}
# elastic_jm_log.txt
2020-06-15 19:52:03.758Z INFO [ult-dispatcher-19107]
o.a.f.r.jobmaster.JobMaster : Trying to recover from a global
failure.org.apache.flink.util.FlinkRuntimeException: Exceeded checkpoint
tolerable failure threshold.2020-06-15 19:52:03.758Z INFO
[ult-dispatcher-19107] o.a.f.r.jobmaster.JobMaster : Trying to recover from
a global failure.org.apache.flink.util.FlinkRuntimeException: Exceeded
checkpoint tolerable failure threshold. at
org.apache.flink.runtime.checkpoint.CheckpointFailureManager.handleTaskLevelCheckpointException(CheckpointFailureManager.java:87)
at
org.apache.flink.runtime.checkpoint.CheckpointCoordinator.failPendingCheckpointDueToTaskFailure(CheckpointCoordinator.java:1467)
at
org.apache.flink.runtime.checkpoint.CheckpointCoordinator.discardCheckpoint(CheckpointCoordinator.java:1377)
at
org.apache.flink.runtime.checkpoint.CheckpointCoordinator.receiveDeclineMessage(CheckpointCoordinator.java:719)
at
org.apache.flink.runtime.scheduler.SchedulerBase.lambda$declineCheckpoint$5(SchedulerBase.java:807)
at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) at
java.util.concurrent.FutureTask.run(FutureTask.java:266) at
java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$201(ScheduledThreadPoolExecutor.java:180)
at
java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:293)
at
java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at
java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:748)
{code}
The job cancellation was stuck with task graph2/graph51/graph52/graph53 for 3
min:
{code:java}
# elastic_tm_log.txt
2020-06-15 19:52:33.841Z WARN [43df85ee0f907ae9d0).] o.a.f.r.taskmanager.Task
: Task 'graph53 (1/1)' did not react to cancelling signal for 30 seconds,
but is stuck in method:
org.elasticsearch.action.bulk.BulkProcessor.flush(BulkProcessor.java:356)
org.apache.flink.streaming.connectors.elasticsearch.ElasticsearchSinkBase.snapshotState(ElasticsearchSinkBase.java:322)
org.apache.flink.streaming.util.functions.StreamingFunctionUtils.trySnapshotFunctionState(StreamingFunctionUtils.java:118)
org.apache.flink.streaming.util.functions.StreamingFunctionUtils.snapshotFunctionState(StreamingFunctionUtils.java:99)
org.apache.flink.streaming.api.operators.AbstractUdfStreamOperator.snapshotState(AbstractUdfStreamOperator.java:90)
org.apache.flink.streaming.api.operators.AbstractStreamOperator.snapshotState(AbstractStreamOperator.java:402)
org.apache.flink.streaming.runtime.tasks.StreamTask$CheckpointingOperation.checkpointStreamOperator(StreamTask.java:1420)
org.apache.flink.streaming.runtime.tasks.StreamTask$CheckpointingOperation.executeCheckpointing(StreamTask.java:1354)
org.apache.flink.streaming.runtime.tasks.StreamTask.checkpointState(StreamTask.java:991)
org.apache.flink.streaming.runtime.tasks.StreamTask.lambda$performCheckpoint$5(StreamTask.java:887)
org.apache.flink.streaming.runtime.tasks.StreamTask$$Lambda$1080/1914359097.run(Unknown
Source)
org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$SynchronizedStreamTaskActionExecutor.runThrowing(StreamTaskActionExecutor.java:94)
org.apache.flink.streaming.runtime.tasks.StreamTask.performCheckpoint(StreamTask.java:860)
org.apache.flink.streaming.runtime.tasks.StreamTask.triggerCheckpoint(StreamTask.java:793)
org.apache.flink.streaming.runtime.tasks.StreamTask.lambda$triggerCheckpointAsync$3(StreamTask.java:777)
org.apache.flink.streaming.runtime.tasks.StreamTask$$Lambda$965/1756450358.call(Unknown
Source)
java.util.concurrent.FutureTask.run(FutureTask.java:266)
org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$SynchronizedStreamTaskActionExecutor.run(StreamTaskActionExecutor.java:87)
org.apache.flink.streaming.runtime.tasks.mailbox.Mail.run(Mail.java:78)
org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.processMail(MailboxProcessor.java:261)
org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:186)
org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:487)
org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:470)
org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:707)
org.apache.flink.runtime.taskmanager.Task.run(Task.java:532)
java.lang.Thread.run(Thread.java:748)2020-06-15 19:52:33.947Z WARN
[5b28276f08a45d7a1e).] o.a.f.r.taskmanager.Task : Task 'graph52 (1/1)'
did not react to cancelling signal for 30 seconds, but is stuck in method:
org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$SynchronizedStreamTaskActionExecutor.run(StreamTaskActionExecutor.java:86)
org.apache.flink.streaming.runtime.tasks.mailbox.Mail.run(Mail.java:78)
org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.processMail(MailboxProcessor.java:261)
org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:186)
org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:487)
org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:470)
org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:707)
org.apache.flink.runtime.taskmanager.Task.run(Task.java:532)
java.lang.Thread.run(Thread.java:748)2020-06-15 19:52:34.104Z WARN
[663038f87ef09c4da6).] o.a.f.r.taskmanager.Task : Task 'graph51 (1/1)'
did not react to cancelling signal for 30 seconds, but is stuck in method:
org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$SynchronizedStreamTaskActionExecutor.run(StreamTaskActionExecutor.java:86)
org.apache.flink.streaming.runtime.tasks.mailbox.Mail.run(Mail.java:78)
org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.processMail(MailboxProcessor.java:261)
org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:186)
org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:487)
org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:470)
org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:707)
org.apache.flink.runtime.taskmanager.Task.run(Task.java:532)
java.lang.Thread.run(Thread.java:748)2020-06-15 19:52:34.162Z WARN
[88fc07dc420a290847).] o.a.f.r.taskmanager.Task : Task 'graph2 (1/1)' did
not react to cancelling signal for 30 seconds, but is stuck in method:
org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$SynchronizedStreamTaskActionExecutor.run(StreamTaskActionExecutor.java:86)
org.apache.flink.streaming.runtime.tasks.mailbox.Mail.run(Mail.java:78)
org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.processMail(MailboxProcessor.java:261)
org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:186)
org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:487)
org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:470)
org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:707)
org.apache.flink.runtime.taskmanager.Task.run(Task.java:532)
java.lang.Thread.run(Thread.java:748)
{code}
Eventually, the TM was shutdown
{code:java}
# elastic_tm_log.txt
2020-06-15 19:55:03.862Z INFO [ault-dispatcher-1717] o.a.f.r.taskmanager.Task
: Attempting to fail task externally graph52 (1/1)
(f6d72e77d5d5e95b28276f08a45d7a1e).
2020-06-15 19:55:03.862Z INFO [ault-dispatcher-1717] o.a.f.r.taskmanager.Task
: Task graph52 (1/1) is already in state CANCELING
2020-06-15 19:55:03.862Z INFO [ault-dispatcher-1717] o.a.f.r.taskmanager.Task
: Attempting to fail task externally graph51 (1/1)
(9ccdf3da055ec3663038f87ef09c4da6).
2020-06-15 19:55:03.862Z INFO [ault-dispatcher-1717] o.a.f.r.taskmanager.Task
: Task graph51 (1/1) is already in state CANCELING
2020-06-15 19:55:03.862Z INFO [ault-dispatcher-1717] o.a.f.r.taskmanager.Task
: Attempting to fail task externally graph2 (1/1)
(0aa096d34c07d288fc07dc420a290847).
2020-06-15 19:55:03.862Z INFO [ault-dispatcher-1717] o.a.f.r.taskmanager.Task
: Task graph2 (1/1) is already in state CANCELING
2020-06-15 19:55:03.862Z INFO [ault-dispatcher-1717] o.a.f.r.taskmanager.Task
: Task graph53 (1/1) is already in state CANCELING
2020-06-15 19:55:03.862Z INFO [ault-dispatcher-1717] o.a.f.r.taskmanager.Task
: Attempting to fail task externally graph53 (1/1)
(e3eb697b20126443df85ee0f907ae9d0).
2020-06-15 19:55:03.870Z INFO [ault-dispatcher-1717]
o.a.f.r.t.s.TaskSlotTableImpl : Free slot TaskSlot(index:0, state:RELEASING,
resource profile: ResourceProfile{cpuCores=1.0000000000000000,
taskHeapMemory=311.040mb (326149077 bytes), taskOffHeapMemory=0 bytes,
managedMemory=42.880mb (44962939 bytes), networkMemory=42.880mb (44962939
bytes)}, allocationId: b059d353dce10b167c9db8bc0d96bce5, jobId:
956d745dc5d43cf6f71df8b00eefd419).
2020-06-15 19:55:03.873Z INFO [ault-dispatcher-1717] o.a.f.r.t.TaskExecutor
: Close JobManager connection for job 956d745dc5d43cf6f71df8b00eefd419.
2020-06-15 19:55:03.941Z ERROR [5b28276f08a45d7a1e).] o.a.f.r.taskmanager.Task
: Task did not exit gracefully within 180 + seconds.
2020-06-15 19:55:03.941Z ERROR [5b28276f08a45d7a1e).]
o.a.f.r.t.TaskManagerRunner : Fatal error occurred while executing the
TaskManager. Shutting it down...
2020-06-15 19:55:03.941Z ERROR [5b28276f08a45d7a1e).] o.a.f.r.t.TaskExecutor
: Task did not exit gracefully within 180 + seconds.
{code}
> ElasticSearch unavailibility causes TM shutdown
> -----------------------------------------------
>
> Key: FLINK-18398
> URL: https://issues.apache.org/jira/browse/FLINK-18398
> Project: Flink
> Issue Type: Bug
> Components: Connectors / ElasticSearch
> Affects Versions: 1.10.0
> Reporter: Alexander Fedulov
> Priority: Critical
> Attachments: elastic_jm_log.txt, elastic_tm_log.txt
>
>
> Similarly to [FLINK-17327|https://issues.apache.org/jira/browse/FLINK-17327],
> unavailibility of ElasticSearch cluster causes Tasks cancellation to timeout
> and Task Manager to be killed. The following exceptions can be found in the
> logs:
>
> {code:java}
> 2020-06-15 19:52:03.664Z ERROR [ I/O dispatcher 229]
> .f.s.c.e.ElasticsearchSinkBase : Failed Elasticsearch bulk request: request
> retries exceeded max retry timeout [30000]java.io.IOException: request
> retries exceeded max retry timeout [30000]
> ...
> 2020-06-15 19:55:03.861Z WARN [43df85ee0f907ae9d0).]
> o.a.f.r.taskmanager.Task : Task 'graph53 (1/1)' did not react to
> cancelling signal for 30 seconds, but is stuck in method:
> org.elasticsearch.action.bulk.BulkProcessor.flush(BulkProcessor.java:356)
> ...
> 2020-06-15 19:55:04.120Z ERROR [663038f87ef09c4da6).]
> o.a.f.r.taskmanager.Task : Task did not exit gracefully within 180 +
> seconds.
> 2020-06-15 19:55:04.121Z ERROR [663038f87ef09c4da6).] o.a.f.r.t.TaskExecutor
> : Task did not exit gracefully within 180 + seconds.
> 2020-06-15 19:55:04.121Z ERROR [663038f87ef09c4da6).]
> o.a.f.r.t.TaskManagerRunner : Fatal error occurred while executing the
> TaskManager. Shutting it down...
> {code}
> Detailed logs are attached.
>
--
This message was sent by Atlassian Jira
(v8.3.4#803005)