This is an automated email from the ASF dual-hosted git repository. dannycranmer pushed a commit to branch master in repository https://gitbox.apache.org/repos/asf/flink.git
commit 16393625ed5e01833ec42e66f533a93c00ac4471 Author: Arvid Heise <[email protected]> AuthorDate: Tue Sep 7 20:03:36 2021 +0200 [FLINK-23528][connectors/kinesis] Gracefully shutdown shard consumer to avoid InterruptionExceptions. --- .../streaming/connectors/kinesis/internals/KinesisDataFetcher.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/flink-connectors/flink-connector-kinesis/src/main/java/org/apache/flink/streaming/connectors/kinesis/internals/KinesisDataFetcher.java b/flink-connectors/flink-connector-kinesis/src/main/java/org/apache/flink/streaming/connectors/kinesis/internals/KinesisDataFetcher.java index 21b3736c789..e12bdaa0b3a 100644 --- a/flink-connectors/flink-connector-kinesis/src/main/java/org/apache/flink/streaming/connectors/kinesis/internals/KinesisDataFetcher.java +++ b/flink-connectors/flink-connector-kinesis/src/main/java/org/apache/flink/streaming/connectors/kinesis/internals/KinesisDataFetcher.java @@ -819,7 +819,7 @@ public class KinesisDataFetcher<T> { LOG.warn("Encountered exception closing record publisher factory", e); } } finally { - shardConsumersExecutor.shutdownNow(); + shardConsumersExecutor.shutdown(); cancelFuture.complete(null);
