[ 
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)

Reply via email to