Hi Arvid,

Thanks for your reply.
In our use case, we are running Flink on EMR cluster.  When we are reading
data from a kinesis stream starting from TRIM_HORIZON offset through an EFO
consumer, the flink source will get these IOException
or StacklessClosedChannelException, which flink will treated as Retryable
exception and enforce the retry limit. And this problem seems mitigated
when we are reading data from LATEST, not as much throughput compared to
reading from TRIM_HORIZON.

We tried to increase the flink parallelism but it doesn't seem to fully fix
the errors. So I am also wondering from a tuning standpoint, what would be
the best approach?
We can certainly raise the retry limit but also 10 is big enough that we
don't want the shard consumption gets stuck there for too long.

Thanks
Leon

On Wed, Dec 15, 2021 at 11:56 AM Arvid Heise <[email protected]> wrote:

> Hi Leon,
>
> The kinesis consumer allows you to specify how many of these exceptions
> can be thrown without a job failure. This is independent from Flink's
> restart policy.
>
> So in your case, if you feel like the Kinesis consumer is giving up too
> quickly, you could simply increase flink.shard.subscribetoshard.maxretries
> to a larger value.
>
> Note that Kinesis distinguishes between RetryableFanOutSubscriberException
> (which you saw and count towards the max retries) and
> RecoverableFanOutSubscriberException (which are repeated indefinitely). So
> the Kinesis connector deems your exceptions as bad enough that they
> probably require a restart at some point in time. At least the second
> exception looks like some network issues; the first one probably is also
> unrecoverable. So if your cluster experiences some network outages from
> time to time, you could just increase the max retries. Of course, that
> means Flink takes longer to restart if something is really off.
>
> If you have more information on the circumstances under which you see
> these exceptions, please share. I'm CCing Danny in case he has more ideas.
>
> On Mon, Dec 13, 2021 at 7:13 AM Leon Xu <[email protected]> wrote:
>
>> Hi Flink users,
>>
>> I used flink-1.12.5 kinesis connector to consume data from kinesis.
>>
>> From time to time I am getting IOException or
>> StacklessClosedChannelException, which will fail the Flink operator when it
>> by default reaches 10 times and trigger the entire job to restart.
>> I have two questions:
>>
>>    1. I am wondering if we can have a better way to handle the error
>>    without restarting the Flink job. Restarting the job is time-consuming and
>>    will often slow down the data consumption.
>>    2. These two errors are treated as retryable exceptions, instead of
>>    recoverable exceptions. Are they not recoverable if we don't restart the
>>    Flink source operator?
>>
>>
>> Followed are the exceptions:
>>
>> org.apache.flink.streaming.connectors.kinesis.internals.publisher.fanout.FanOutShardSubscriber$RetryableFanOutSubscriberException:
>> org.apache.flink.kinesis.shaded.io.netty.channel.StacklessClosedChannelException
>> at
>> org.apache.flink.streaming.connectors.kinesis.internals.publisher.fanout.FanOutShardSubscriber.handleError(FanOutShardSubscriber.java:296)
>> ~[blob_p-6581ced9ade704dee24c2632e701b413d2e656ba-f7515b4b21d25cb44b07411a9ad294e8:?]
>> at
>> org.apache.flink.streaming.connectors.kinesis.internals.publisher.fanout.FanOutShardSubscriber.consumeAllRecordsFromKinesisShard(FanOutShardSubscriber.java:363)
>> ~[blob_p-6581ced9ade704dee24c2632e701b413d2e656ba-f7515b4b21d25cb44b07411a9ad294e8:?]
>> at
>> org.apache.flink.streaming.connectors.kinesis.internals.publisher.fanout.FanOutShardSubscriber.subscribeToShardAndConsumeRecords(FanOutShardSubscriber.java:188)
>> ~[blob_p-6581ced9ade704dee24c2632e701b413d2e656ba-f7515b4b21d25cb44b07411a9ad294e8:?]
>> at
>> org.apache.flink.streaming.connectors.kinesis.internals.publisher.fanout.FanOutRecordPublisher.runWithBackoff(FanOutRecordPublisher.java:154)
>> [blob_p-6581ced9ade704dee24c2632e701b413d2e656ba-f7515b4b21d25cb44b07411a9ad294e8:?]
>>
>> or
>>
>> org.apache.flink.streaming.connectors.kinesis.internals.publisher.fanout.FanOutShardSubscriber$RetryableFanOutSubscriberException:
>> java.io.IOException: An error occurred on the connection: null at
>> org.apache.flink.streaming.connectors.kinesis.internals.publisher.fanout.FanOutShardSubscriber.handleError(FanOutShardSubscriber.java:296)
>> ~[blob_p-6581ced9ade704dee24c2632e701b413d2e656ba-f7515b4b21d25cb44b07411a9ad294e8:?]
>> at
>> org.apache.flink.streaming.connectors.kinesis.internals.publisher.fanout.FanOutShardSubscriber.consumeAllRecordsFromKinesisShard(FanOutShardSubscriber.java:363)
>> ~[blob_p-6581ced9ade704dee24c2632e701b413d2e656ba-f7515b4b21d25cb44b07411a9ad294e8:?]
>> at
>> org.apache.flink.streaming.connectors.kinesis.internals.publisher.fanout.FanOutShardSubscriber.subscribeToShardAndConsumeRecords(FanOutShardSubscriber.java:188)
>> ~[blob_p-6581ced9ade704dee24c2632e701b413d2e656ba-f7515b4b21d25cb44b07411a9ad294e8:?]
>> at
>> org.apache.flink.streaming.connectors.kinesis.internals.publisher.fanout.FanOutRecordPublisher.runWithBackoff(FanOutRecordPublisher.java:154)
>> [blob_p-6581ced9ade704dee24c2632e701b413d2e656ba-f7515b4b21d25cb44b07411a9ad294e8:?]
>>
>>
>>
>> Thanks
>>
>> Leon
>>
>>
>>

Reply via email to