This is an automated email from the ASF dual-hosted git repository.

tison pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-connector-pulsar.git


The following commit(s) were added to refs/heads/main by this push:
     new 2278653  [FLINK-30109][Connector/Pulsar] Drop the use of sneaky 
exception. (#24)
2278653 is described below

commit 2278653d67a8ddf171c88d538a288e503221625a
Author: Yufan Sheng <[email protected]>
AuthorDate: Mon Feb 13 21:46:26 2023 +0800

    [FLINK-30109][Connector/Pulsar] Drop the use of sneaky exception. (#24)
---
 .../pulsar/common/config/PulsarClientFactory.java  |  21 ++--
 .../pulsar/common/utils/PulsarExceptionUtils.java  |  83 ------------
 .../common/utils/PulsarTransactionUtils.java       |  26 ++--
 .../flink/connector/pulsar/sink/PulsarSink.java    |   5 +-
 .../pulsar/sink/committer/PulsarCommitter.java     |   7 +-
 .../pulsar/sink/config/SinkConfiguration.java      |   2 +-
 .../connector/pulsar/sink/writer/PulsarWriter.java |   6 +-
 .../pulsar/sink/writer/topic/MetadataListener.java |   4 +-
 .../pulsar/sink/writer/topic/ProducerRegister.java |  31 +++--
 .../connector/pulsar/source/PulsarSource.java      |   7 +-
 .../pulsar/source/PulsarSourceBuilder.java         |   2 +-
 .../source/enumerator/PulsarSourceEnumerator.java  |  53 ++++----
 .../source/enumerator/cursor/StopCursor.java       |   2 +-
 .../cursor/stop/LatestMessageStopCursor.java       |   7 +-
 .../enumerator/subscriber/PulsarSubscriber.java    |  15 ++-
 .../subscriber/impl/BasePulsarSubscriber.java      |  50 +++++---
 .../subscriber/impl/TopicListSubscriber.java       |  32 ++---
 .../subscriber/impl/TopicPatternSubscriber.java    |  74 ++++++-----
 .../source/reader/PulsarPartitionSplitReader.java  |  29 +++--
 .../source/reader/PulsarSourceFetcherManager.java  |  22 ++--
 .../pulsar/source/split/PulsarPartitionSplit.java  |   7 +-
 .../pulsar/common/schema/PulsarSchemaTest.java     |   1 +
 .../sink/writer/topic/MetadataListenerTest.java    |   6 +-
 .../sink/writer/topic/ProducerRegisterTest.java    |   6 +-
 .../enumerator/PulsarSourceEnumeratorTest.java     |  25 ++--
 .../source/enumerator/cursor/StopCursorTest.java   |   4 +-
 .../subscriber/PulsarSubscriberTest.java           |  41 +++---
 .../reader/PulsarPartitionSplitReaderTest.java     |  74 +++++------
 .../source/reader/PulsarSourceReaderTest.java      |   2 +-
 .../PulsarDeserializationSchemaTest.java           |  16 +--
 .../pulsar/testutils/PulsarTestCommonUtils.java    |   6 -
 .../pulsar/testutils/PulsarTestEnvironment.java    |   8 +-
 .../pulsar/testutils/function/ControlSource.java   |  14 +--
 .../pulsar/testutils/runtime/PulsarRuntime.java    |   4 +-
 .../testutils/runtime/PulsarRuntimeOperator.java   | 140 ++++++++-------------
 .../runtime/container/PulsarContainerRuntime.java  |  16 +--
 .../testutils/sink/PulsarSinkTestContext.java      |   7 +-
 .../sink/reader/PulsarPartitionDataReader.java     |  10 +-
 .../cases/MultipleTopicsConsumingContext.java      |  12 +-
 .../source/cases/SingleTopicConsumingContext.java  |  17 ++-
 .../writer/KeyedPulsarPartitionDataWriter.java     |  15 ++-
 .../source/writer/PulsarEncryptDataWriter.java     |  19 ++-
 .../source/writer/PulsarPartitionDataWriter.java   |   7 +-
 43 files changed, 453 insertions(+), 482 deletions(-)

diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/common/config/PulsarClientFactory.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/common/config/PulsarClientFactory.java
index 1f01b24..4939a23 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/common/config/PulsarClientFactory.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/common/config/PulsarClientFactory.java
@@ -27,6 +27,7 @@ import org.apache.pulsar.client.api.AuthenticationFactory;
 import org.apache.pulsar.client.api.ClientBuilder;
 import org.apache.pulsar.client.api.ProxyProtocol;
 import org.apache.pulsar.client.api.PulsarClient;
+import org.apache.pulsar.client.api.PulsarClientException;
 import org.apache.pulsar.client.impl.auth.AuthenticationDisabled;
 
 import java.util.Map;
@@ -76,7 +77,6 @@ import static 
org.apache.flink.connector.pulsar.common.config.PulsarOptions.PULS
 import static 
org.apache.flink.connector.pulsar.common.config.PulsarOptions.PULSAR_TLS_TRUST_STORE_TYPE;
 import static 
org.apache.flink.connector.pulsar.common.config.PulsarOptions.PULSAR_USE_KEY_STORE_TLS;
 import static 
org.apache.flink.connector.pulsar.common.config.PulsarOptions.PULSAR_USE_TCP_NO_DELAY;
-import static 
org.apache.flink.connector.pulsar.common.utils.PulsarExceptionUtils.sneakyClient;
 import static org.apache.pulsar.client.api.SizeUnit.BYTES;
 
 /** The factory for creating pulsar client classes from {@link 
PulsarConfiguration}. */
@@ -88,7 +88,8 @@ public final class PulsarClientFactory {
     }
 
     /** Create a PulsarClient by using the flink Configuration and the config 
customizer. */
-    public static PulsarClient createClient(PulsarConfiguration configuration) 
{
+    public static PulsarClient createClient(PulsarConfiguration configuration)
+            throws PulsarClientException {
         ClientBuilder builder = PulsarClient.builder();
 
         // requestTimeoutMs don't have a setter method on ClientBuilder. We 
have to use low level
@@ -148,14 +149,15 @@ public final class PulsarClientFactory {
         }
         configuration.useOption(PULSAR_ENABLE_TRANSACTION, 
builder::enableTransaction);
 
-        return sneakyClient(builder::build);
+        return builder.build();
     }
 
     /**
      * PulsarAdmin shares almost the same configuration with PulsarClient, but 
we separate this
      * creating method for directly use it.
      */
-    public static PulsarAdmin createAdmin(PulsarConfiguration configuration) {
+    public static PulsarAdmin createAdmin(PulsarConfiguration configuration)
+            throws PulsarClientException {
         PulsarAdminBuilder builder = PulsarAdmin.builder();
 
         // Create the authentication instance for the Pulsar client.
@@ -182,7 +184,7 @@ public final class PulsarClientFactory {
         configuration.useOption(
                 PULSAR_AUTO_CERT_REFRESH_TIME, v -> 
builder.autoCertRefreshTime(v, MILLISECONDS));
 
-        return sneakyClient(builder::build);
+        return builder.build();
     }
 
     /**
@@ -192,14 +194,14 @@ public final class PulsarClientFactory {
      *
      * <p>This method behavior is the same as the pulsar command line tools.
      */
-    private static Authentication createAuthentication(PulsarConfiguration 
configuration) {
+    private static Authentication createAuthentication(PulsarConfiguration 
configuration)
+            throws PulsarClientException {
         if (configuration.contains(PULSAR_AUTH_PLUGIN_CLASS_NAME)) {
             String authPluginClassName = 
configuration.get(PULSAR_AUTH_PLUGIN_CLASS_NAME);
 
             if (configuration.contains(PULSAR_AUTH_PARAMS)) {
                 String authParamsString = 
configuration.get(PULSAR_AUTH_PARAMS);
-                return sneakyClient(
-                        () -> 
AuthenticationFactory.create(authPluginClassName, authParamsString));
+                return AuthenticationFactory.create(authPluginClassName, 
authParamsString);
             } else {
                 Map<String, String> paramsMap = 
configuration.getProperties(PULSAR_AUTH_PARAM_MAP);
                 if (paramsMap.isEmpty()) {
@@ -209,8 +211,7 @@ public final class PulsarClientFactory {
                                     PULSAR_AUTH_PARAMS.key(), 
PULSAR_AUTH_PARAM_MAP.key()));
                 }
 
-                return sneakyClient(
-                        () -> 
AuthenticationFactory.create(authPluginClassName, paramsMap));
+                return AuthenticationFactory.create(authPluginClassName, 
paramsMap);
             }
         }
 
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/common/utils/PulsarExceptionUtils.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/common/utils/PulsarExceptionUtils.java
deleted file mode 100644
index 63bd65c..0000000
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/common/utils/PulsarExceptionUtils.java
+++ /dev/null
@@ -1,83 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements.  See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership.  The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.flink.connector.pulsar.common.utils;
-
-import org.apache.flink.annotation.Internal;
-import org.apache.flink.util.function.SupplierWithException;
-import org.apache.flink.util.function.ThrowingRunnable;
-
-import org.apache.pulsar.client.admin.PulsarAdminException;
-import org.apache.pulsar.client.api.PulsarClientException;
-
-/**
- * Util class for pulsar checked exceptions. Sneaky throw {@link 
PulsarAdminException} and {@link
- * PulsarClientException}.
- */
-@Internal
-public final class PulsarExceptionUtils {
-
-    private PulsarExceptionUtils() {
-        // No public constructor.
-    }
-
-    public static <R extends PulsarClientException> void sneakyClient(
-            ThrowingRunnable<R> runnable) {
-        sneaky(runnable);
-    }
-
-    public static <T, R extends PulsarClientException> T sneakyClient(
-            SupplierWithException<T, R> supplier) {
-        return sneaky(supplier);
-    }
-
-    public static <R extends PulsarAdminException> void 
sneakyAdmin(ThrowingRunnable<R> runnable) {
-        sneaky(runnable);
-    }
-
-    public static <T, R extends PulsarAdminException> T sneakyAdmin(
-            SupplierWithException<T, R> supplier) {
-        return sneaky(supplier);
-    }
-
-    private static <R extends Exception> void sneaky(ThrowingRunnable<R> 
runnable) {
-        try {
-            runnable.run();
-        } catch (Exception r) {
-            sneakyThrow(r);
-        }
-    }
-
-    /** Catch the throwable exception and rethrow it without try catch. */
-    private static <T, R extends Exception> T sneaky(SupplierWithException<T, 
R> supplier) {
-        try {
-            return supplier.get();
-        } catch (Exception r) {
-            sneakyThrow(r);
-        }
-
-        // This method wouldn't be executed.
-        throw new RuntimeException("Never throw here.");
-    }
-
-    /** javac hack for unchecking the checked exception. */
-    @SuppressWarnings("unchecked")
-    public static <T extends Exception> void sneakyThrow(Exception t) throws T 
{
-        throw (T) t;
-    }
-}
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/common/utils/PulsarTransactionUtils.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/common/utils/PulsarTransactionUtils.java
index fbb597d..e377320 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/common/utils/PulsarTransactionUtils.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/common/utils/PulsarTransactionUtils.java
@@ -19,21 +19,17 @@
 package org.apache.flink.connector.pulsar.common.utils;
 
 import org.apache.flink.annotation.Internal;
-import org.apache.flink.util.FlinkRuntimeException;
 
 import org.apache.pulsar.client.api.PulsarClient;
+import org.apache.pulsar.client.api.PulsarClientException;
 import org.apache.pulsar.client.api.transaction.Transaction;
 import org.apache.pulsar.client.api.transaction.TransactionCoordinatorClient;
 import org.apache.pulsar.client.api.transaction.TxnID;
 import org.apache.pulsar.client.impl.PulsarClientImpl;
 
-import java.util.concurrent.CompletableFuture;
-import java.util.concurrent.ExecutionException;
 import java.util.concurrent.TimeUnit;
 
-import static 
org.apache.flink.connector.pulsar.common.utils.PulsarExceptionUtils.sneakyClient;
 import static org.apache.flink.util.Preconditions.checkNotNull;
-import static 
org.apache.pulsar.client.api.transaction.TransactionCoordinatorClientException.unwrap;
 
 /** A suit of workarounds for the Pulsar Transaction. */
 @Internal
@@ -44,19 +40,19 @@ public final class PulsarTransactionUtils {
     }
 
     /** Create transaction with given timeout millis. */
-    public static Transaction createTransaction(PulsarClient pulsarClient, 
long timeoutMs) {
+    public static Transaction createTransaction(PulsarClient pulsarClient, 
long timeoutMs)
+            throws PulsarClientException {
         try {
-            CompletableFuture<Transaction> future =
-                    sneakyClient(pulsarClient::newTransaction)
-                            .withTransactionTimeout(timeoutMs, 
TimeUnit.MILLISECONDS)
-                            .build();
-
-            return future.get();
+            return pulsarClient
+                    .newTransaction()
+                    .withTransactionTimeout(timeoutMs, TimeUnit.MILLISECONDS)
+                    .build()
+                    .get();
         } catch (InterruptedException e) {
             Thread.currentThread().interrupt();
-            throw new IllegalStateException(e);
-        } catch (ExecutionException e) {
-            throw new FlinkRuntimeException(unwrap(e));
+            throw new PulsarClientException(e);
+        } catch (Exception e) {
+            throw PulsarClientException.unwrap(e);
         }
     }
 
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/PulsarSink.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/PulsarSink.java
index 883c8ea..b7daaae 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/PulsarSink.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/PulsarSink.java
@@ -38,6 +38,8 @@ import 
org.apache.flink.connector.pulsar.sink.writer.serializer.PulsarSerializat
 import org.apache.flink.connector.pulsar.sink.writer.topic.MetadataListener;
 import org.apache.flink.core.io.SimpleVersionedSerializer;
 
+import org.apache.pulsar.client.api.PulsarClientException;
+
 import javax.annotation.Nullable;
 
 import static org.apache.flink.util.Preconditions.checkNotNull;
@@ -128,7 +130,8 @@ public class PulsarSink<IN> implements 
TwoPhaseCommittingSink<IN, PulsarCommitta
 
     @Internal
     @Override
-    public PrecommittingSinkWriter<IN, PulsarCommittable> 
createWriter(InitContext initContext) {
+    public PrecommittingSinkWriter<IN, PulsarCommittable> 
createWriter(InitContext initContext)
+            throws PulsarClientException {
         return new PulsarWriter<>(
                 sinkConfiguration,
                 serializationSchema,
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/committer/PulsarCommitter.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/committer/PulsarCommitter.java
index 11e80c7..7058409 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/committer/PulsarCommitter.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/committer/PulsarCommitter.java
@@ -26,6 +26,7 @@ import 
org.apache.flink.connector.pulsar.sink.config.SinkConfiguration;
 import org.apache.flink.util.FlinkRuntimeException;
 
 import org.apache.pulsar.client.api.PulsarClient;
+import org.apache.pulsar.client.api.PulsarClientException;
 import org.apache.pulsar.client.api.transaction.TransactionCoordinatorClient;
 import 
org.apache.pulsar.client.api.transaction.TransactionCoordinatorClientException;
 import 
org.apache.pulsar.client.api.transaction.TransactionCoordinatorClientException.CoordinatorNotFoundException;
@@ -65,7 +66,8 @@ public class PulsarCommitter implements 
Committer<PulsarCommittable>, Closeable
 
     @Override
     @SuppressWarnings("java:S3776")
-    public void commit(Collection<CommitRequest<PulsarCommittable>> requests) {
+    public void commit(Collection<CommitRequest<PulsarCommittable>> requests)
+            throws PulsarClientException {
         TransactionCoordinatorClient client = transactionCoordinatorClient();
 
         for (CommitRequest<PulsarCommittable> request : requests) {
@@ -147,7 +149,8 @@ public class PulsarCommitter implements 
Committer<PulsarCommittable>, Closeable
      * DeliveryGuarantee#NONE} and {@link DeliveryGuarantee#AT_LEAST_ONCE}. So 
we couldn't create
      * the Pulsar client at first.
      */
-    private TransactionCoordinatorClient transactionCoordinatorClient() {
+    private TransactionCoordinatorClient transactionCoordinatorClient()
+            throws PulsarClientException {
         if (coordinatorClient == null) {
             this.pulsarClient = createClient(sinkConfiguration);
             this.coordinatorClient = getTcClient(pulsarClient);
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/config/SinkConfiguration.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/config/SinkConfiguration.java
index d041bde..e5d6dc5 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/config/SinkConfiguration.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/config/SinkConfiguration.java
@@ -19,7 +19,7 @@
 package org.apache.flink.connector.pulsar.sink.config;
 
 import org.apache.flink.annotation.PublicEvolving;
-import org.apache.flink.api.connector.sink.Sink.InitContext;
+import org.apache.flink.api.connector.sink2.Sink.InitContext;
 import org.apache.flink.configuration.Configuration;
 import org.apache.flink.connector.base.DeliveryGuarantee;
 import org.apache.flink.connector.pulsar.common.config.PulsarConfiguration;
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/writer/PulsarWriter.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/writer/PulsarWriter.java
index 0cd51f7..cc47609 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/writer/PulsarWriter.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/writer/PulsarWriter.java
@@ -43,6 +43,7 @@ import org.apache.flink.util.FlinkRuntimeException;
 import org.apache.flink.shaded.guava30.com.google.common.base.Strings;
 
 import org.apache.pulsar.client.api.MessageId;
+import org.apache.pulsar.client.api.PulsarClientException;
 import org.apache.pulsar.client.api.Schema;
 import org.apache.pulsar.client.api.TypedMessageBuilder;
 import org.slf4j.Logger;
@@ -100,7 +101,8 @@ public class PulsarWriter<IN> implements 
PrecommittingSinkWriter<IN, PulsarCommi
             TopicRouter<IN> topicRouter,
             MessageDelayer<IN> messageDelayer,
             PulsarCrypto pulsarCrypto,
-            InitContext initContext) {
+            InitContext initContext)
+            throws PulsarClientException {
         checkNotNull(sinkConfiguration);
         this.serializationSchema = checkNotNull(serializationSchema);
         this.metadataListener = checkNotNull(metadataListener);
@@ -183,7 +185,7 @@ public class PulsarWriter<IN> implements 
PrecommittingSinkWriter<IN, PulsarCommi
 
     @SuppressWarnings({"rawtypes", "unchecked"})
     private TypedMessageBuilder<?> createMessageBuilder(
-            String topic, Context context, PulsarMessage<?> message) {
+            String topic, Context context, PulsarMessage<?> message) throws 
PulsarClientException {
 
         Schema<?> schema = message.getSchema();
         TypedMessageBuilder<?> builder = 
producerRegister.createMessageBuilder(topic, schema);
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/writer/topic/MetadataListener.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/writer/topic/MetadataListener.java
index bbf909c..796a60e 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/writer/topic/MetadataListener.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/writer/topic/MetadataListener.java
@@ -36,6 +36,7 @@ import 
org.apache.flink.shaded.guava30.com.google.common.collect.ImmutableList;
 import org.apache.pulsar.client.admin.PulsarAdmin;
 import org.apache.pulsar.client.admin.PulsarAdminException;
 import org.apache.pulsar.client.admin.PulsarAdminException.NotFoundException;
+import org.apache.pulsar.client.api.PulsarClientException;
 import org.apache.pulsar.common.naming.TopicName;
 import org.apache.pulsar.common.partition.PartitionedTopicMetadata;
 import org.slf4j.Logger;
@@ -99,7 +100,8 @@ public class MetadataListener implements Serializable, 
Closeable {
     }
 
     /** Register the topic metadata update action in process time service. */
-    public void open(SinkConfiguration sinkConfiguration, 
ProcessingTimeService timeService) {
+    public void open(SinkConfiguration sinkConfiguration, 
ProcessingTimeService timeService)
+            throws PulsarClientException {
         // Initialize listener properties.
         this.pulsarAdmin = createAdmin(sinkConfiguration);
         this.topicMetadataRefreshInterval = 
sinkConfiguration.getTopicMetadataRefreshInterval();
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/writer/topic/ProducerRegister.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/writer/topic/ProducerRegister.java
index f57faca..08d4eb9 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/writer/topic/ProducerRegister.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/sink/writer/topic/ProducerRegister.java
@@ -37,6 +37,7 @@ import org.apache.pulsar.client.api.Producer;
 import org.apache.pulsar.client.api.ProducerBuilder;
 import org.apache.pulsar.client.api.ProducerStats;
 import org.apache.pulsar.client.api.PulsarClient;
+import org.apache.pulsar.client.api.PulsarClientException;
 import org.apache.pulsar.client.api.Schema;
 import org.apache.pulsar.client.api.TypedMessageBuilder;
 import org.apache.pulsar.client.api.transaction.Transaction;
@@ -85,7 +86,6 @@ import static 
org.apache.flink.connector.pulsar.common.metrics.MetricNames.TOTAL
 import static 
org.apache.flink.connector.pulsar.common.metrics.MetricNames.TOTAL_BYTES_SENT;
 import static 
org.apache.flink.connector.pulsar.common.metrics.MetricNames.TOTAL_MSGS_SENT;
 import static 
org.apache.flink.connector.pulsar.common.metrics.MetricNames.TOTAL_SEND_FAILED;
-import static 
org.apache.flink.connector.pulsar.common.utils.PulsarExceptionUtils.sneakyClient;
 import static 
org.apache.flink.connector.pulsar.common.utils.PulsarTransactionUtils.createTransaction;
 import static 
org.apache.flink.connector.pulsar.common.utils.PulsarTransactionUtils.getTcClient;
 import static 
org.apache.flink.connector.pulsar.sink.config.PulsarSinkConfigUtils.createProducerBuilder;
@@ -112,7 +112,8 @@ public class ProducerRegister implements Closeable {
     public ProducerRegister(
             SinkConfiguration sinkConfiguration,
             PulsarCrypto pulsarCrypto,
-            SinkWriterMetricGroup metricGroup) {
+            SinkWriterMetricGroup metricGroup)
+            throws PulsarClientException {
         this.pulsarClient = createClient(sinkConfiguration);
         this.sinkConfiguration = sinkConfiguration;
         this.pulsarCrypto = pulsarCrypto;
@@ -143,8 +144,8 @@ public class ProducerRegister implements Closeable {
      * have to manually create it.
      */
     @SuppressWarnings("unchecked")
-    public <T> TypedMessageBuilder<T> createMessageBuilder(
-            String topic, @Nullable Schema<?> schema) {
+    public <T> TypedMessageBuilder<T> createMessageBuilder(String topic, 
@Nullable Schema<?> schema)
+            throws PulsarClientException {
         if (schema == null || schema.getSchemaInfo().getType() == 
SchemaType.BYTES) {
             schema = getBytesSchema(topic);
         }
@@ -209,7 +210,8 @@ public class ProducerRegister implements Closeable {
 
     /** Create or return the cached topic-related producer. */
     @SuppressWarnings("unchecked")
-    private <T> Producer<T> getOrCreateProducer(String topic, Schema<T> 
schema) {
+    private <T> Producer<T> getOrCreateProducer(String topic, Schema<T> schema)
+            throws PulsarClientException {
         Map<SchemaHash, Producer<?>> set = producers.computeIfAbsent(topic, t 
-> new HashMap<>());
         SchemaHash hash = PulsarSchemaUtils.hash(schema);
         if (set.containsKey(hash)) {
@@ -241,7 +243,7 @@ public class ProducerRegister implements Closeable {
         // Set the sending counter for metrics.
         builder.intercept(new ProducerMetricsInterceptor(metricGroup));
 
-        Producer<T> producer = sneakyClient(builder::create);
+        Producer<T> producer = builder.create();
 
         // Expose the stats for calculating and monitoring.
         exposeProducerMetrics(producer);
@@ -280,13 +282,16 @@ public class ProducerRegister implements Closeable {
     /**
      * Get the cached topic-related transaction. Or create a new transaction 
after checkpointing.
      */
-    private Transaction getOrCreateTransaction(String topic) {
-        return transactions.computeIfAbsent(
-                topic,
-                t -> {
-                    long timeoutMillis = 
sinkConfiguration.getTransactionTimeoutMillis();
-                    return createTransaction(pulsarClient, timeoutMillis);
-                });
+    private Transaction getOrCreateTransaction(String topic) throws 
PulsarClientException {
+        if (transactions.containsKey(topic)) {
+            return transactions.get(topic);
+        }
+
+        long timeoutMillis = sinkConfiguration.getTransactionTimeoutMillis();
+        Transaction transaction = createTransaction(pulsarClient, 
timeoutMillis);
+        transactions.put(topic, transaction);
+
+        return transaction;
     }
 
     /**
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/PulsarSource.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/PulsarSource.java
index 5dc9d94..c7d272b 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/PulsarSource.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/PulsarSource.java
@@ -43,6 +43,8 @@ import 
org.apache.flink.connector.pulsar.source.split.PulsarPartitionSplit;
 import 
org.apache.flink.connector.pulsar.source.split.PulsarPartitionSplitSerializer;
 import org.apache.flink.core.io.SimpleVersionedSerializer;
 
+import org.apache.pulsar.client.api.PulsarClientException;
+
 /**
  * The Source implementation of Pulsar. Please use a {@link 
PulsarSourceBuilder} to construct a
  * {@link PulsarSource}. The following example shows how to create a 
PulsarSource emitting records
@@ -139,7 +141,7 @@ public final class PulsarSource<OUT>
     @Internal
     @Override
     public SplitEnumerator<PulsarPartitionSplit, PulsarSourceEnumState> 
createEnumerator(
-            SplitEnumeratorContext<PulsarPartitionSplit> enumContext) {
+            SplitEnumeratorContext<PulsarPartitionSplit> enumContext) throws 
PulsarClientException {
         return new PulsarSourceEnumerator(
                 subscriber,
                 startCursor,
@@ -153,7 +155,8 @@ public final class PulsarSource<OUT>
     @Override
     public SplitEnumerator<PulsarPartitionSplit, PulsarSourceEnumState> 
restoreEnumerator(
             SplitEnumeratorContext<PulsarPartitionSplit> enumContext,
-            PulsarSourceEnumState checkpoint) {
+            PulsarSourceEnumState checkpoint)
+            throws PulsarClientException {
         return new PulsarSourceEnumerator(
                 subscriber,
                 startCursor,
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/PulsarSourceBuilder.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/PulsarSourceBuilder.java
index d15c0e3..80d8c30 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/PulsarSourceBuilder.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/PulsarSourceBuilder.java
@@ -619,7 +619,7 @@ public final class PulsarSourceBuilder<OUT> {
 
     private void ensureSchemaTypeIsValid(Schema<?> schema) {
         SchemaInfo info = schema.getSchemaInfo();
-        if (info.getType() == SchemaType.AUTO_CONSUME || info.getType() == 
SchemaType.AUTO) {
+        if (info.getType() == SchemaType.AUTO_CONSUME) {
             throw new IllegalArgumentException(
                     "Auto schema is only supported by providing a 
GenericRecordDeserializer");
         }
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/PulsarSourceEnumerator.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/PulsarSourceEnumerator.java
index 4287565..82a15a9 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/PulsarSourceEnumerator.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/PulsarSourceEnumerator.java
@@ -34,6 +34,8 @@ import 
org.apache.flink.metrics.groups.SplitEnumeratorMetricGroup;
 import org.apache.flink.util.FlinkRuntimeException;
 
 import org.apache.pulsar.client.admin.PulsarAdmin;
+import org.apache.pulsar.client.admin.PulsarAdminException;
+import org.apache.pulsar.client.api.PulsarClientException;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -45,7 +47,6 @@ import java.util.Set;
 
 import static java.util.Collections.singletonList;
 import static 
org.apache.flink.connector.pulsar.common.config.PulsarClientFactory.createAdmin;
-import static 
org.apache.flink.connector.pulsar.common.utils.PulsarExceptionUtils.sneakyAdmin;
 import static 
org.apache.flink.connector.pulsar.source.enumerator.PulsarSourceEnumState.initialState;
 import static 
org.apache.flink.connector.pulsar.source.enumerator.assigner.SplitAssigner.createAssigner;
 
@@ -71,7 +72,8 @@ public class PulsarSourceEnumerator
             StopCursor stopCursor,
             RangeGenerator rangeGenerator,
             SourceConfiguration sourceConfiguration,
-            SplitEnumeratorContext<PulsarPartitionSplit> context) {
+            SplitEnumeratorContext<PulsarPartitionSplit> context)
+            throws PulsarClientException {
         this(
                 subscriber,
                 startCursor,
@@ -89,7 +91,8 @@ public class PulsarSourceEnumerator
             RangeGenerator rangeGenerator,
             SourceConfiguration sourceConfiguration,
             SplitEnumeratorContext<PulsarPartitionSplit> context,
-            PulsarSourceEnumState enumState) {
+            PulsarSourceEnumState enumState)
+            throws PulsarClientException {
         this.pulsarAdmin = createAdmin(sourceConfiguration);
         this.subscriber = subscriber;
         this.startCursor = startCursor;
@@ -102,6 +105,7 @@ public class PulsarSourceEnumerator
 
     @Override
     public void start() {
+        subscriber.open(pulsarAdmin);
         rangeGenerator.open(sourceConfiguration);
 
         // Expose the split assignment metrics if Flink has supported.
@@ -166,7 +170,7 @@ public class PulsarSourceEnumerator
     }
 
     @Override
-    public void close() {
+    public void close() throws PulsarClientException {
         if (pulsarAdmin != null) {
             pulsarAdmin.close();
         }
@@ -182,9 +186,9 @@ public class PulsarSourceEnumerator
      *
      * @return Set of subscribed {@link TopicPartition}s
      */
-    private Set<TopicPartition> getSubscribedTopicPartitions() {
+    private Set<TopicPartition> getSubscribedTopicPartitions() throws 
Exception {
         int parallelism = context.currentParallelism();
-        return subscriber.getSubscribedTopicPartitions(pulsarAdmin, 
rangeGenerator, parallelism);
+        return subscriber.getSubscribedTopicPartitions(rangeGenerator, 
parallelism);
     }
 
     /**
@@ -199,38 +203,45 @@ public class PulsarSourceEnumerator
     private void checkPartitionChanges(Set<TopicPartition> fetchedPartitions, 
Throwable throwable) {
         if (throwable != null) {
             throw new FlinkRuntimeException(
-                    "Failed to list subscribed topic partitions due to ", 
throwable);
+                    "Failed to list subscribed topic partitions due to: " + 
throwable.getMessage(),
+                    throwable);
         }
 
-        // Append the partitions into current assignment state.
+        // Append the partitions into current assignment state twice,
+        // because the getSubscribedTopicPartitions method is executed in 
another thread.
         List<TopicPartition> newPartitions =
                 splitAssigner.registerTopicPartitions(fetchedPartitions);
-        createSubscription(newPartitions);
 
-        // Assign the new readers.
-        List<Integer> registeredReaders = new 
ArrayList<>(context.registeredReaders().keySet());
-        assignPendingPartitionSplits(registeredReaders);
-    }
-
-    /** Create subscription on topic partition if it doesn't exist. */
-    private void createSubscription(List<TopicPartition> newPartitions) {
+        // Create subscription on newly discovered topic partitions if it 
doesn't contain related
+        // subscription.
         for (TopicPartition partition : newPartitions) {
             String topic = partition.getFullTopicName();
             String subscriptionName = 
sourceConfiguration.getSubscriptionName();
             CursorPosition position =
                     startCursor.position(partition.getTopic(), 
partition.getPartitionId());
 
-            if (sourceConfiguration.isResetSubscriptionCursor()) {
-                sneakyAdmin(() -> position.seekPosition(pulsarAdmin, topic, 
subscriptionName));
-            } else {
-                sneakyAdmin(
-                        () -> position.createInitialPosition(pulsarAdmin, 
topic, subscriptionName));
+            try {
+                if (sourceConfiguration.isResetSubscriptionCursor()) {
+                    position.seekPosition(pulsarAdmin, topic, 
subscriptionName);
+                } else {
+                    position.createInitialPosition(pulsarAdmin, topic, 
subscriptionName);
+                }
+            } catch (PulsarAdminException e) {
+                throw new FlinkRuntimeException(e);
             }
         }
+
+        // Assign the new readers.
+        List<Integer> registeredReaders = new 
ArrayList<>(context.registeredReaders().keySet());
+        assignPendingPartitionSplits(registeredReaders);
     }
 
     /** Query the unassigned splits and assign them to the available readers. 
*/
     private void assignPendingPartitionSplits(List<Integer> pendingReaders) {
+        if (pendingReaders.isEmpty()) {
+            return;
+        }
+
         // Validate the reader.
         pendingReaders.forEach(
                 reader -> {
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/cursor/StopCursor.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/cursor/StopCursor.java
index 866d036..1af875e 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/cursor/StopCursor.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/cursor/StopCursor.java
@@ -42,7 +42,7 @@ import java.io.Serializable;
 public interface StopCursor extends Serializable {
 
     /** The open method for the cursor initializer. This method could be 
executed multiple times. */
-    default void open(PulsarAdmin admin, TopicPartition partition) {}
+    default void open(PulsarAdmin admin, TopicPartition partition) throws 
Exception {}
 
     /** Determine whether to pause consumption on the current message by the 
returned enum. */
     StopCondition shouldStop(Message<?> message);
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/cursor/stop/LatestMessageStopCursor.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/cursor/stop/LatestMessageStopCursor.java
index 0de963e..5af8be8 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/cursor/stop/LatestMessageStopCursor.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/cursor/stop/LatestMessageStopCursor.java
@@ -22,11 +22,10 @@ import 
org.apache.flink.connector.pulsar.source.enumerator.cursor.StopCursor;
 import 
org.apache.flink.connector.pulsar.source.enumerator.topic.TopicPartition;
 
 import org.apache.pulsar.client.admin.PulsarAdmin;
+import org.apache.pulsar.client.admin.PulsarAdminException;
 import org.apache.pulsar.client.api.Message;
 import org.apache.pulsar.client.api.MessageId;
 
-import static 
org.apache.flink.connector.pulsar.common.utils.PulsarExceptionUtils.sneakyAdmin;
-
 /**
  * A stop cursor that initialize the position to the latest message id. The 
offsets initialization
  * are taken care of by the {@code PulsarPartitionSplitReaderBase} instead of 
by the {@code
@@ -49,10 +48,10 @@ public class LatestMessageStopCursor implements StopCursor {
     }
 
     @Override
-    public void open(PulsarAdmin admin, TopicPartition partition) {
+    public void open(PulsarAdmin admin, TopicPartition partition) throws 
PulsarAdminException {
         if (messageId == null) {
             String topic = partition.getFullTopicName();
-            this.messageId = sneakyAdmin(() -> 
admin.topics().getLastMessageId(topic));
+            this.messageId = admin.topics().getLastMessageId(topic);
         }
     }
 }
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/PulsarSubscriber.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/PulsarSubscriber.java
index b8a55bf..f4102e9 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/PulsarSubscriber.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/PulsarSubscriber.java
@@ -50,13 +50,20 @@ public interface PulsarSubscriber extends Serializable {
      * Get a set of subscribed {@link TopicPartition}s. The method could throw 
{@link
      * IllegalStateException}, an extra try catch is required.
      *
-     * @param pulsarAdmin The admin interface used to retrieve subscribed 
topic partitions.
-     * @param rangeGenerator The range for different partitions.
+     * @param generator The range for different partitions.
      * @param parallelism The parallelism of flink source.
      * @return A subscribed {@link TopicPartition} for each pulsar topic 
partition.
      */
-    Set<TopicPartition> getSubscribedTopicPartitions(
-            PulsarAdmin pulsarAdmin, RangeGenerator rangeGenerator, int 
parallelism);
+    @SuppressWarnings("java:S112")
+    Set<TopicPartition> getSubscribedTopicPartitions(RangeGenerator generator, 
int parallelism)
+            throws Exception;
+
+    /**
+     * Initialize the topic subscriber.
+     *
+     * @param admin The admin interface used to retrieve subscribed topic 
partitions.
+     */
+    void open(PulsarAdmin admin);
 
     // ----------------- factory methods --------------
 
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/impl/BasePulsarSubscriber.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/impl/BasePulsarSubscriber.java
index a206a25..306646c 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/impl/BasePulsarSubscriber.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/impl/BasePulsarSubscriber.java
@@ -23,26 +23,28 @@ import 
org.apache.flink.connector.pulsar.source.enumerator.topic.TopicMetadata;
 import 
org.apache.flink.connector.pulsar.source.enumerator.topic.TopicNameUtils;
 import 
org.apache.flink.connector.pulsar.source.enumerator.topic.TopicPartition;
 import org.apache.flink.connector.pulsar.source.enumerator.topic.TopicRange;
+import 
org.apache.flink.connector.pulsar.source.enumerator.topic.range.RangeGenerator;
 
 import org.apache.pulsar.client.admin.PulsarAdmin;
 import org.apache.pulsar.client.admin.PulsarAdminException;
 import org.apache.pulsar.common.partition.PartitionedTopicMetadata;
 
-import java.util.ArrayList;
+import java.util.HashSet;
 import java.util.List;
-
-import static java.util.Collections.singletonList;
+import java.util.Set;
 
 /** PulsarSubscriber abstract class to simplify Pulsar admin related 
operations. */
 public abstract class BasePulsarSubscriber implements PulsarSubscriber {
     private static final long serialVersionUID = 2053021503331058888L;
 
-    protected TopicMetadata queryTopicMetadata(PulsarAdmin pulsarAdmin, String 
topicName) {
+    protected transient PulsarAdmin admin;
+
+    protected TopicMetadata queryTopicMetadata(String topicName) throws 
PulsarAdminException {
         // Drop the complete topic name for a clean partitioned topic name.
         String completeTopicName = TopicNameUtils.topicName(topicName);
         try {
             PartitionedTopicMetadata metadata =
-                    
pulsarAdmin.topics().getPartitionedTopicMetadata(completeTopicName);
+                    
admin.topics().getPartitionedTopicMetadata(completeTopicName);
             return new TopicMetadata(topicName, metadata.partitions);
         } catch (PulsarAdminException e) {
             if (e.getStatusCode() == 404) {
@@ -50,23 +52,37 @@ public abstract class BasePulsarSubscriber implements 
PulsarSubscriber {
                 return null;
             } else {
                 // This method would cause failure for subscribers.
-                throw new IllegalStateException(e);
+                throw e;
             }
         }
     }
 
-    protected List<TopicPartition> toTopicPartitions(
-            TopicMetadata metadata, List<TopicRange> ranges) {
-        if (!metadata.isPartitioned()) {
-            // For non-partitioned topic.
-            return singletonList(new TopicPartition(metadata.getName(), 
ranges));
-        } else {
-            // For partitioned topic.
-            List<TopicPartition> partitions = new ArrayList<>();
-            for (int i = 0; i < metadata.getPartitionSize(); i++) {
-                partitions.add(new TopicPartition(metadata.getName(), i, 
ranges));
+    protected Set<TopicPartition> createTopicPartitions(
+            Set<String> topics, RangeGenerator generator, int parallelism)
+            throws PulsarAdminException {
+        Set<TopicPartition> results = new HashSet<>();
+
+        for (String topic : topics) {
+            TopicMetadata metadata = queryTopicMetadata(topic);
+            if (metadata != null) {
+                List<TopicRange> ranges = generator.range(metadata, 
parallelism);
+                if (!metadata.isPartitioned()) {
+                    // For non-partitioned topic.
+                    results.add(new TopicPartition(metadata.getName(), 
ranges));
+                } else {
+                    // For partitioned topic.
+                    for (int i = 0; i < metadata.getPartitionSize(); i++) {
+                        results.add(new TopicPartition(metadata.getName(), i, 
ranges));
+                    }
+                }
             }
-            return partitions;
         }
+
+        return results;
+    }
+
+    @Override
+    public void open(PulsarAdmin admin) {
+        this.admin = admin;
     }
 }
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/impl/TopicListSubscriber.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/impl/TopicListSubscriber.java
index 1b87cc2..fe0863a 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/impl/TopicListSubscriber.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/impl/TopicListSubscriber.java
@@ -23,10 +23,9 @@ import 
org.apache.flink.connector.pulsar.source.enumerator.topic.TopicPartition;
 import org.apache.flink.connector.pulsar.source.enumerator.topic.TopicRange;
 import 
org.apache.flink.connector.pulsar.source.enumerator.topic.range.RangeGenerator;
 
-import org.apache.pulsar.client.admin.PulsarAdmin;
+import org.apache.pulsar.client.admin.PulsarAdminException;
 import org.apache.pulsar.common.naming.TopicName;
 
-import java.util.ArrayList;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Set;
@@ -37,12 +36,12 @@ import static 
org.apache.flink.connector.pulsar.source.enumerator.topic.TopicNam
 public class TopicListSubscriber extends BasePulsarSubscriber {
     private static final long serialVersionUID = 6473918213832993116L;
 
-    private final List<String> partitions;
-    private final List<String> fullTopicNames;
+    private final Set<String> partitions;
+    private final Set<String> fullTopicNames;
 
     public TopicListSubscriber(List<String> fullTopicNameOrPartitions) {
-        this.partitions = new ArrayList<>();
-        this.fullTopicNames = new ArrayList<>();
+        this.partitions = new HashSet<>();
+        this.fullTopicNames = new HashSet<>();
 
         for (String fullTopicNameOrPartition : fullTopicNameOrPartitions) {
             if (isPartition(fullTopicNameOrPartition)) {
@@ -55,26 +54,21 @@ public class TopicListSubscriber extends 
BasePulsarSubscriber {
 
     @Override
     public Set<TopicPartition> getSubscribedTopicPartitions(
-            PulsarAdmin pulsarAdmin, RangeGenerator rangeGenerator, int 
parallelism) {
-        Set<TopicPartition> results = new HashSet<>();
+            RangeGenerator generator, int parallelism) throws 
PulsarAdminException {
 
         // Query topics from Pulsar.
-        for (String topic : fullTopicNames) {
-            TopicMetadata metadata = queryTopicMetadata(pulsarAdmin, topic);
-            List<TopicRange> ranges = rangeGenerator.range(metadata, 
parallelism);
-
-            results.addAll(toTopicPartitions(metadata, ranges));
-        }
+        Set<TopicPartition> results = createTopicPartitions(fullTopicNames, 
generator, parallelism);
 
+        // Query partitions from Pulsar.
         for (String partition : partitions) {
             TopicName topicName = TopicName.get(partition);
             String name = topicName.getPartitionedTopicName();
             int index = topicName.getPartitionIndex();
-
-            TopicMetadata metadata = queryTopicMetadata(pulsarAdmin, name);
-            List<TopicRange> ranges = rangeGenerator.range(metadata, 
parallelism);
-
-            results.add(new TopicPartition(name, index, ranges));
+            TopicMetadata metadata = queryTopicMetadata(name);
+            if (metadata != null) {
+                List<TopicRange> ranges = generator.range(metadata, 
parallelism);
+                results.add(new TopicPartition(name, index, ranges));
+            }
         }
 
         return results;
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/impl/TopicPatternSubscriber.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/impl/TopicPatternSubscriber.java
index 146996e..dc77df0 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/impl/TopicPatternSubscriber.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/impl/TopicPatternSubscriber.java
@@ -18,74 +18,68 @@
 
 package org.apache.flink.connector.pulsar.source.enumerator.subscriber.impl;
 
-import 
org.apache.flink.connector.pulsar.source.enumerator.topic.TopicNameUtils;
 import 
org.apache.flink.connector.pulsar.source.enumerator.topic.TopicPartition;
-import org.apache.flink.connector.pulsar.source.enumerator.topic.TopicRange;
 import 
org.apache.flink.connector.pulsar.source.enumerator.topic.range.RangeGenerator;
 
-import org.apache.pulsar.client.admin.PulsarAdmin;
-import org.apache.pulsar.client.admin.PulsarAdminException;
 import org.apache.pulsar.client.api.RegexSubscriptionMode;
 import org.apache.pulsar.common.naming.NamespaceName;
 import org.apache.pulsar.common.naming.TopicName;
 
-import java.util.Collections;
+import java.util.HashSet;
 import java.util.List;
-import java.util.Objects;
 import java.util.Set;
 import java.util.regex.Pattern;
 
-import static java.util.stream.Collectors.toSet;
-import static 
org.apache.flink.shaded.guava30.com.google.common.base.Predicates.not;
+import static 
org.apache.flink.connector.pulsar.source.enumerator.topic.TopicNameUtils.isInternal;
+import static 
org.apache.flink.connector.pulsar.source.enumerator.topic.TopicNameUtils.topicName;
 
 /** Subscribe to matching topics based on topic pattern. */
 public class TopicPatternSubscriber extends BasePulsarSubscriber {
     private static final long serialVersionUID = 3307710093243745104L;
 
-    private final Pattern topicPattern;
-    private final RegexSubscriptionMode subscriptionMode;
+    private final Pattern shortenedPattern;
     private final String namespace;
+    private final RegexSubscriptionMode subscriptionMode;
 
     public TopicPatternSubscriber(Pattern topicPattern, RegexSubscriptionMode 
subscriptionMode) {
-        this.topicPattern = topicPattern;
         this.subscriptionMode = subscriptionMode;
+        String pattern = topicPattern.toString();
+        this.shortenedPattern =
+                pattern.contains("://") ? 
Pattern.compile(pattern.split("://")[1]) : topicPattern;
 
         // Extract the namespace from topic pattern regex.
         // If no namespace provided in the regex, we would directly use 
"default" as the namespace.
-        TopicName destination = TopicName.get(topicPattern.toString());
+        TopicName destination = TopicName.get(topicPattern.pattern());
         NamespaceName namespaceName = destination.getNamespaceObject();
         this.namespace = namespaceName.toString();
     }
 
     @Override
     public Set<TopicPartition> getSubscribedTopicPartitions(
-            PulsarAdmin pulsarAdmin, RangeGenerator rangeGenerator, int 
parallelism) {
-        try {
-            return pulsarAdmin
-                    .namespaces()
-                    .getTopics(namespace)
-                    .parallelStream()
-                    .filter(this::matchesSubscriptionMode)
-                    .filter(not(TopicNameUtils::isInternal))
-                    .filter(topic -> topicPattern.matcher(topic).find())
-                    .map(topic -> queryTopicMetadata(pulsarAdmin, topic))
-                    .filter(Objects::nonNull)
-                    .flatMap(
-                            metadata -> {
-                                List<TopicRange> ranges =
-                                        rangeGenerator.range(metadata, 
parallelism);
-                                return toTopicPartitions(metadata, 
ranges).stream();
-                            })
-                    .collect(toSet());
-        } catch (PulsarAdminException e) {
-            if (e.getStatusCode() == 404) {
-                // Skip the topic metadata query.
-                return Collections.emptySet();
-            } else {
-                // This method would cause failure for subscribers.
-                throw new IllegalStateException(e);
+            RangeGenerator generator, int parallelism) throws Exception {
+        // This method will query a set of existed topic partitions.
+        List<String> partitions = admin.namespaces().getTopics(namespace);
+        Set<String> results = new HashSet<>(partitions.size());
+
+        for (String partition : partitions) {
+            String topic = topicName(partition);
+            if (matchesSubscriptionMode(topic)
+                    && !isInternal(topic)
+                    && matchesTopicPattern(topic)) {
+                results.add(topic);
             }
         }
+
+        return createTopicPartitions(results, generator, parallelism);
+    }
+
+    /**
+     * If the topic matches 'topicsPattern'. This method is in the 
PulsarClient, and it's removed
+     * since 2.11.0 release. We keep the method here.
+     */
+    private boolean matchesTopicPattern(String topic) {
+        String shortenedTopic = topic.split("://")[1];
+        return shortenedPattern.matcher(shortenedTopic).matches();
     }
 
     /**
@@ -100,9 +94,11 @@ public class TopicPatternSubscriber extends 
BasePulsarSubscriber {
                 return topicName.isPersistent();
             case NonPersistentOnly:
                 return !topicName.isPersistent();
-            default:
-                // RegexSubscriptionMode.AllTopics
+            case AllTopics:
                 return true;
+            default:
+                throw new IllegalArgumentException(
+                        "We don't support such subscription mode " + 
subscriptionMode);
         }
     }
 }
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/reader/PulsarPartitionSplitReader.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/reader/PulsarPartitionSplitReader.java
index 6ba3274..ca4dbc3 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/reader/PulsarPartitionSplitReader.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/reader/PulsarPartitionSplitReader.java
@@ -35,6 +35,7 @@ import 
org.apache.flink.connector.pulsar.source.enumerator.topic.TopicPartition;
 import org.apache.flink.connector.pulsar.source.split.PulsarPartitionSplit;
 import org.apache.flink.metrics.MetricGroup;
 import org.apache.flink.metrics.groups.SourceReaderMetricGroup;
+import org.apache.flink.util.FlinkRuntimeException;
 import org.apache.flink.util.Preconditions;
 
 import org.apache.flink.shaded.guava30.com.google.common.base.Strings;
@@ -50,6 +51,7 @@ import org.apache.pulsar.client.api.Message;
 import org.apache.pulsar.client.api.MessageCrypto;
 import org.apache.pulsar.client.api.MessageId;
 import org.apache.pulsar.client.api.PulsarClient;
+import org.apache.pulsar.client.api.PulsarClientException;
 import org.apache.pulsar.client.api.Schema;
 import org.apache.pulsar.common.api.proto.MessageMetadata;
 import org.slf4j.Logger;
@@ -78,7 +80,6 @@ import static 
org.apache.flink.connector.pulsar.common.metrics.MetricNames.TOTAL
 import static 
org.apache.flink.connector.pulsar.common.metrics.MetricNames.TOTAL_BYTES_RECEIVED;
 import static 
org.apache.flink.connector.pulsar.common.metrics.MetricNames.TOTAL_MSGS_RECEIVED;
 import static 
org.apache.flink.connector.pulsar.common.metrics.MetricNames.TOTAL_RECEIVED_FAILED;
-import static 
org.apache.flink.connector.pulsar.common.utils.PulsarExceptionUtils.sneakyClient;
 import static 
org.apache.flink.connector.pulsar.source.config.CursorVerification.FAIL_ON_MISMATCH;
 import static 
org.apache.flink.connector.pulsar.source.config.PulsarSourceConfigUtils.createConsumerBuilder;
 import static 
org.apache.flink.connector.pulsar.source.enumerator.topic.range.TopicRangeUtils.isFullTopicRanges;
@@ -191,7 +192,11 @@ public class PulsarPartitionSplitReader
         this.registeredSplit = newSplits.get(0);
 
         // Open stop cursor.
-        registeredSplit.open(pulsarAdmin);
+        try {
+            registeredSplit.open(pulsarAdmin);
+        } catch (Exception e) {
+            throw new FlinkRuntimeException(e);
+        }
 
         // Reset the start position before creating the consumer.
         MessageId latestConsumedId = registeredSplit.getLatestConsumedId();
@@ -231,7 +236,11 @@ public class PulsarPartitionSplitReader
         }
 
         // Create pulsar consumer.
-        this.pulsarConsumer = 
createPulsarConsumer(registeredSplit.getPartition());
+        try {
+            this.pulsarConsumer = 
createPulsarConsumer(registeredSplit.getPartition());
+        } catch (PulsarClientException e) {
+            throw new FlinkRuntimeException(e);
+        }
 
         LOG.info("Register split {} consumer for current reader.", 
registeredSplit);
     }
@@ -258,24 +267,26 @@ public class PulsarPartitionSplitReader
     }
 
     @Override
-    public void close() {
+    public void close() throws PulsarClientException {
         if (pulsarConsumer != null) {
-            sneakyClient(() -> pulsarConsumer.close());
+            pulsarConsumer.close();
         }
     }
 
-    public void notifyCheckpointComplete(TopicPartition partition, MessageId 
offsetsToCommit) {
+    public void notifyCheckpointComplete(TopicPartition partition, MessageId 
offsetsToCommit)
+            throws PulsarClientException {
         if (pulsarConsumer == null) {
             this.pulsarConsumer = createPulsarConsumer(partition);
         }
 
-        sneakyClient(() -> 
pulsarConsumer.acknowledgeCumulative(offsetsToCommit));
+        pulsarConsumer.acknowledgeCumulative(offsetsToCommit);
     }
 
     // --------------------------- Helper Methods -----------------------------
 
     /** Create a specified {@link Consumer} by the given topic partition. */
-    private Consumer<byte[]> createPulsarConsumer(TopicPartition partition) {
+    private Consumer<byte[]> createPulsarConsumer(TopicPartition partition)
+            throws PulsarClientException {
         ConsumerBuilder<byte[]> consumerBuilder =
                 createConsumerBuilder(pulsarClient, schema, 
sourceConfiguration);
 
@@ -304,7 +315,7 @@ public class PulsarPartitionSplitReader
         }
 
         // Create the consumer configuration by using common utils.
-        Consumer<byte[]> consumer = sneakyClient(consumerBuilder::subscribe);
+        Consumer<byte[]> consumer = consumerBuilder.subscribe();
 
         // Exposing the consumer metrics.
         exposeConsumerMetrics(consumer);
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/reader/PulsarSourceFetcherManager.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/reader/PulsarSourceFetcherManager.java
index 8d8e802..7382f29 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/reader/PulsarSourceFetcherManager.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/reader/PulsarSourceFetcherManager.java
@@ -32,6 +32,7 @@ import 
org.apache.flink.connector.pulsar.source.split.PulsarPartitionSplit;
 import org.apache.pulsar.client.api.Consumer;
 import org.apache.pulsar.client.api.Message;
 import org.apache.pulsar.client.api.MessageId;
+import org.apache.pulsar.client.api.PulsarClientException;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -105,20 +106,25 @@ public class PulsarSourceFetcherManager
         }
     }
 
-    public void acknowledgeMessages(Map<TopicPartition, MessageId> 
cursorsToCommit) {
+    public void acknowledgeMessages(Map<TopicPartition, MessageId> 
cursorsToCommit)
+            throws PulsarClientException {
         LOG.debug("Acknowledge messages {}", cursorsToCommit);
-        cursorsToCommit.forEach(
-                (partition, messageId) -> {
-                    SplitFetcher<Message<byte[]>, PulsarPartitionSplit> 
fetcher =
-                            getOrCreateFetcher(partition.toString());
-                    triggerAcknowledge(fetcher, partition, messageId);
-                });
+
+        for (Map.Entry<TopicPartition, MessageId> entry : 
cursorsToCommit.entrySet()) {
+            TopicPartition partition = entry.getKey();
+            MessageId messageId = entry.getValue();
+
+            SplitFetcher<Message<byte[]>, PulsarPartitionSplit> fetcher =
+                    getOrCreateFetcher(partition.toString());
+            triggerAcknowledge(fetcher, partition, messageId);
+        }
     }
 
     private void triggerAcknowledge(
             SplitFetcher<Message<byte[]>, PulsarPartitionSplit> splitFetcher,
             TopicPartition partition,
-            MessageId messageId) {
+            MessageId messageId)
+            throws PulsarClientException {
         PulsarPartitionSplitReader splitReader =
                 (PulsarPartitionSplitReader) splitFetcher.getSplitReader();
         splitReader.notifyCheckpointComplete(partition, messageId);
diff --git 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/split/PulsarPartitionSplit.java
 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/split/PulsarPartitionSplit.java
index 3189c17..3044e9c 100644
--- 
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/split/PulsarPartitionSplit.java
+++ 
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/split/PulsarPartitionSplit.java
@@ -22,6 +22,7 @@ import org.apache.flink.annotation.Internal;
 import org.apache.flink.api.connector.source.SourceSplit;
 import org.apache.flink.connector.pulsar.source.enumerator.cursor.StopCursor;
 import 
org.apache.flink.connector.pulsar.source.enumerator.topic.TopicPartition;
+import org.apache.flink.connector.pulsar.source.reader.PulsarSourceReader;
 
 import org.apache.pulsar.client.admin.PulsarAdmin;
 import org.apache.pulsar.client.api.MessageId;
@@ -44,8 +45,8 @@ public class PulsarPartitionSplit implements SourceSplit, 
Serializable {
     private final StopCursor stopCursor;
 
     /**
-     * Since this field in only used in {@link 
PulsarOrderedSourceReader#snapshotState(long)}, it's
-     * no need to serialize this field into flink checkpoint state.
+     * Since this field in only used in {@link 
PulsarSourceReader#snapshotState(long)}, it's no need
+     * to serialize this field into flink checkpoint state.
      */
     @Nullable private final MessageId latestConsumedId;
 
@@ -93,7 +94,7 @@ public class PulsarPartitionSplit implements SourceSplit, 
Serializable {
     }
 
     /** Open stop cursor. */
-    public void open(PulsarAdmin admin) {
+    public void open(PulsarAdmin admin) throws Exception {
         stopCursor.open(admin, partition);
     }
 
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/common/schema/PulsarSchemaTest.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/common/schema/PulsarSchemaTest.java
index 81074c9..9547812 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/common/schema/PulsarSchemaTest.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/common/schema/PulsarSchemaTest.java
@@ -85,6 +85,7 @@ class PulsarSchemaTest {
     }
 
     @Test
+    @SuppressWarnings("unchecked")
     void invalidPulsarSchemaCreationWithoutClassType() {
         assertThrows(IllegalArgumentException.class, () -> new 
PulsarSchema<>(AVRO));
         assertThrows(IllegalArgumentException.class, () -> new 
PulsarSchema<>(JSON));
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/sink/writer/topic/MetadataListenerTest.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/sink/writer/topic/MetadataListenerTest.java
index 998a115..8485fae 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/sink/writer/topic/MetadataListenerTest.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/sink/writer/topic/MetadataListenerTest.java
@@ -42,7 +42,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals;
 class MetadataListenerTest extends PulsarTestSuiteBase {
 
     @Test
-    void listenEmptyTopics() {
+    void listenEmptyTopics() throws Exception {
         MetadataListener listener = new MetadataListener();
         SinkConfiguration configuration = 
sinkConfiguration(Duration.ofMinutes(5).toMillis());
         TestProcessingTimeService timeService = new 
TestProcessingTimeService();
@@ -80,7 +80,7 @@ class MetadataListenerTest extends PulsarTestSuiteBase {
     }
 
     @Test
-    void fetchTopicPartitionInformation() {
+    void fetchTopicPartitionInformation() throws Exception {
         String topic = randomAlphabetic(10);
         operator().createTopic(topic, 8);
 
@@ -128,7 +128,7 @@ class MetadataListenerTest extends PulsarTestSuiteBase {
     }
 
     @Test
-    void fetchNonPartitionTopic() {
+    void fetchNonPartitionTopic() throws Exception {
         String topic = randomAlphabetic(10);
         operator().createTopic(topic, 0);
         List<TopicPartition> nonPartitionTopic = singletonList(new 
TopicPartition(topic));
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/sink/writer/topic/ProducerRegisterTest.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/sink/writer/topic/ProducerRegisterTest.java
index 75b5703..8329e8b 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/sink/writer/topic/ProducerRegisterTest.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/sink/writer/topic/ProducerRegisterTest.java
@@ -35,7 +35,6 @@ import org.junit.jupiter.api.Test;
 import org.junit.jupiter.params.ParameterizedTest;
 import org.junit.jupiter.params.provider.EnumSource;
 
-import java.io.IOException;
 import java.util.List;
 import java.util.concurrent.ThreadLocalRandom;
 
@@ -53,7 +52,7 @@ class ProducerRegisterTest extends PulsarTestSuiteBase {
     @ParameterizedTest
     @EnumSource(DeliveryGuarantee.class)
     void createMessageBuilderForSendingMessage(DeliveryGuarantee 
deliveryGuarantee)
-            throws IOException {
+            throws Exception {
         String topic = randomAlphabetic(10);
         operator().createTopic(topic, 8);
 
@@ -83,7 +82,8 @@ class ProducerRegisterTest extends PulsarTestSuiteBase {
     @EnumSource(
             value = DeliveryGuarantee.class,
             names = {"AT_LEAST_ONCE", "NONE"})
-    void noneAndAtLeastOnceWouldNotCreateTransaction(DeliveryGuarantee 
deliveryGuarantee) {
+    void noneAndAtLeastOnceWouldNotCreateTransaction(DeliveryGuarantee 
deliveryGuarantee)
+            throws Exception {
         String topic = randomAlphabetic(10);
         operator().createTopic(topic, 8);
 
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/enumerator/PulsarSourceEnumeratorTest.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/enumerator/PulsarSourceEnumeratorTest.java
index f79357d..dd1cf20 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/enumerator/PulsarSourceEnumeratorTest.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/enumerator/PulsarSourceEnumeratorTest.java
@@ -44,6 +44,7 @@ import java.util.Set;
 import java.util.concurrent.TimeUnit;
 import java.util.regex.Pattern;
 
+import static java.util.stream.Collectors.joining;
 import static java.util.stream.Collectors.toSet;
 import static org.apache.commons.lang3.RandomStringUtils.randomAlphabetic;
 import static 
org.apache.flink.connector.pulsar.source.PulsarSourceOptions.PULSAR_PARTITION_DISCOVERY_INTERVAL_MS;
@@ -57,6 +58,7 @@ import static org.assertj.core.api.Assertions.assertThat;
 /** Unit tests for {@link PulsarSourceEnumerator}. */
 class PulsarSourceEnumeratorTest extends PulsarTestSuiteBase {
 
+    private static final String TOPIC_PREFIX = "enumerator-topic-";
     private static final int NUM_SUBTASKS = 3;
     private static final int READER0 = 0;
     private static final int READER1 = 1;
@@ -126,7 +128,7 @@ class PulsarSourceEnumeratorTest extends 
PulsarTestSuiteBase {
 
     @Test
     void discoverPartitionsPeriodically() throws Throwable {
-        String dynamicTopic = "topic3-" + randomAlphabetic(10);
+        String dynamicTopic = TOPIC_PREFIX + randomAlphabetic(10);
         Set<String> preexistingTopics = setupPreexistingTopics();
         Set<String> topicsToSubscribe = new HashSet<>(preexistingTopics);
         topicsToSubscribe.add(dynamicTopic);
@@ -241,9 +243,9 @@ class PulsarSourceEnumeratorTest extends 
PulsarTestSuiteBase {
         }
     }
 
-    private Set<String> setupPreexistingTopics() {
-        String topic1 = "topic1-" + randomAlphabetic(10);
-        String topic2 = "topic2-" + randomAlphabetic(10);
+    private Set<String> setupPreexistingTopics() throws Exception {
+        String topic1 = "enumerator-topic-" + randomAlphabetic(10);
+        String topic2 = "enumerator-topic-" + randomAlphabetic(10);
 
         operator().setupTopic(topic1);
         operator().setupTopic(topic2);
@@ -274,7 +276,8 @@ class PulsarSourceEnumeratorTest extends 
PulsarTestSuiteBase {
     private PulsarSourceEnumerator createEnumerator(
             Set<String> topics,
             MockSplitEnumeratorContext<PulsarPartitionSplit> enumContext,
-            boolean enablePeriodicPartitionDiscovery) {
+            boolean enablePeriodicPartitionDiscovery)
+            throws Exception {
         return createEnumerator(
                 topics, enumContext, enablePeriodicPartitionDiscovery, 
initialState());
     }
@@ -283,10 +286,18 @@ class PulsarSourceEnumeratorTest extends 
PulsarTestSuiteBase {
             Set<String> topicsToSubscribe,
             MockSplitEnumeratorContext<PulsarPartitionSplit> enumContext,
             boolean enablePeriodicPartitionDiscovery,
-            PulsarSourceEnumState sourceEnumState) {
+            PulsarSourceEnumState sourceEnumState)
+            throws Exception {
         // Use a TopicPatternSubscriber so that no exception if a subscribed 
topic hasn't been
         // created yet.
-        String topicRegex = String.join("|", topicsToSubscribe);
+        String topicRegex =
+                topicsToSubscribe.stream()
+                        .map(s -> s.substring(TOPIC_PREFIX.length()))
+                        .collect(
+                                joining(
+                                        "|",
+                                        "persistent://public/default/" + 
TOPIC_PREFIX + "(",
+                                        ")"));
         Pattern topicPattern = Pattern.compile(topicRegex);
         PulsarSubscriber subscriber =
                 getTopicPatternSubscriber(topicPattern, 
RegexSubscriptionMode.AllTopics);
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/enumerator/cursor/StopCursorTest.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/enumerator/cursor/StopCursorTest.java
index e09180f..7cd685d 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/enumerator/cursor/StopCursorTest.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/enumerator/cursor/StopCursorTest.java
@@ -33,8 +33,6 @@ import org.apache.pulsar.client.api.MessageId;
 import org.apache.pulsar.client.api.Schema;
 import org.junit.jupiter.api.Test;
 
-import java.io.IOException;
-
 import static java.util.Collections.singletonList;
 import static org.apache.commons.lang3.RandomStringUtils.randomAlphabetic;
 import static org.apache.commons.lang3.RandomStringUtils.randomAlphanumeric;
@@ -51,7 +49,7 @@ import static org.assertj.core.api.Assertions.assertThat;
 class StopCursorTest extends PulsarTestSuiteBase {
 
     @Test
-    void publishTimeStopCursor() throws IOException {
+    void publishTimeStopCursor() throws Exception {
         String topicName = "stop-cursor-" + randomAlphanumeric(5);
         operator().createTopic(topicName, 2);
 
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/PulsarSubscriberTest.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/PulsarSubscriberTest.java
index bedc3d8..b8cc415 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/PulsarSubscriberTest.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/enumerator/subscriber/PulsarSubscriberTest.java
@@ -58,7 +58,7 @@ class PulsarSubscriberTest extends PulsarTestSuiteBase {
     private static final int NUM_PARALLELISM = 10;
 
     @BeforeAll
-    void setUp() {
+    void setUp() throws Exception {
         operator().createTopic(topic1, NUM_PARTITIONS_PER_TOPIC);
         operator().createTopic(topic2, NUM_PARTITIONS_PER_TOPIC);
         operator().createTopic(topic3, NUM_PARTITIONS_PER_TOPIC);
@@ -67,7 +67,7 @@ class PulsarSubscriberTest extends PulsarTestSuiteBase {
     }
 
     @AfterAll
-    void tearDown() {
+    void tearDown() throws Exception {
         operator().deleteTopic(topic1);
         operator().deleteTopic(topic2);
         operator().deleteTopic(topic3);
@@ -76,11 +76,12 @@ class PulsarSubscriberTest extends PulsarTestSuiteBase {
     }
 
     @Test
-    void topicListSubscriber() {
+    void topicListSubscriber() throws Exception {
         PulsarSubscriber subscriber = 
getTopicListSubscriber(Arrays.asList(topic1, topic2));
+        subscriber.open(operator().admin());
+
         Set<TopicPartition> topicPartitions =
-                subscriber.getSubscribedTopicPartitions(
-                        operator().admin(), new FullRangeGenerator(), 
NUM_PARALLELISM);
+                subscriber.getSubscribedTopicPartitions(new 
FullRangeGenerator(), NUM_PARALLELISM);
         Set<TopicPartition> expectedPartitions = new HashSet<>();
 
         for (int i = 0; i < NUM_PARTITIONS_PER_TOPIC; i++) {
@@ -92,40 +93,42 @@ class PulsarSubscriberTest extends PulsarTestSuiteBase {
     }
 
     @Test
-    void subscribeOnePartitionOfMultiplePartitionTopic() {
+    void subscribeOnePartitionOfMultiplePartitionTopic() throws Exception {
         String partition = topicNameWithPartition(topic1, 2);
 
         PulsarSubscriber subscriber = 
getTopicListSubscriber(singletonList(partition));
+        subscriber.open(operator().admin());
+
         Set<TopicPartition> partitions =
-                subscriber.getSubscribedTopicPartitions(
-                        operator().admin(), new FullRangeGenerator(), 
NUM_PARALLELISM);
+                subscriber.getSubscribedTopicPartitions(new 
FullRangeGenerator(), NUM_PARALLELISM);
 
         TopicPartition desiredPartition = new TopicPartition(topic1, 2);
         assertThat(partitions).hasSize(1).containsExactly(desiredPartition);
     }
 
     @Test
-    void subscribeNonPartitionedTopicList() {
+    void subscribeNonPartitionedTopicList() throws Exception {
         PulsarSubscriber subscriber = 
getTopicListSubscriber(singletonList(topic4));
+        subscriber.open(operator().admin());
+
         Set<TopicPartition> partitions =
-                subscriber.getSubscribedTopicPartitions(
-                        operator().admin(), new FullRangeGenerator(), 
NUM_PARALLELISM);
+                subscriber.getSubscribedTopicPartitions(new 
FullRangeGenerator(), NUM_PARALLELISM);
 
         TopicPartition desiredPartition = new TopicPartition(topic4);
         assertThat(partitions).hasSize(1).containsExactly(desiredPartition);
     }
 
     @Test
-    void subscribeNonPartitionedTopicPattern() {
+    void subscribeNonPartitionedTopicPattern() throws Exception {
         PulsarSubscriber subscriber =
                 getTopicPatternSubscriber(
                         Pattern.compile(
-                                
"persistent://public/default/pulsar-subscriber-non-partitioned-topic*?"),
+                                
"persistent://public/default/pulsar-subscriber-non-partitioned-topic.*?"),
                         AllTopics);
+        subscriber.open(operator().admin());
 
         Set<TopicPartition> topicPartitions =
-                subscriber.getSubscribedTopicPartitions(
-                        operator().admin(), new FullRangeGenerator(), 
NUM_PARALLELISM);
+                subscriber.getSubscribedTopicPartitions(new 
FullRangeGenerator(), NUM_PARALLELISM);
 
         Set<TopicPartition> expectedPartitions = new HashSet<>();
 
@@ -136,15 +139,15 @@ class PulsarSubscriberTest extends PulsarTestSuiteBase {
     }
 
     @Test
-    void topicPatternSubscriber() {
+    void topicPatternSubscriber() throws Exception {
         PulsarSubscriber subscriber =
                 getTopicPatternSubscriber(
-                        
Pattern.compile("persistent://public/default/pulsar-subscriber-topic*?"),
+                        
Pattern.compile("persistent://public/default/pulsar-subscriber-topic.*?"),
                         AllTopics);
+        subscriber.open(operator().admin());
 
         Set<TopicPartition> topicPartitions =
-                subscriber.getSubscribedTopicPartitions(
-                        operator().admin(), new FullRangeGenerator(), 
NUM_PARALLELISM);
+                subscriber.getSubscribedTopicPartitions(new 
FullRangeGenerator(), NUM_PARALLELISM);
 
         Set<TopicPartition> expectedPartitions = new HashSet<>();
 
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/reader/PulsarPartitionSplitReaderTest.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/reader/PulsarPartitionSplitReaderTest.java
index 9d5f65b..0fdced1 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/reader/PulsarPartitionSplitReaderTest.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/reader/PulsarPartitionSplitReaderTest.java
@@ -40,13 +40,11 @@ import java.io.IOException;
 import java.util.ArrayList;
 import java.util.List;
 import java.util.concurrent.TimeUnit;
-import java.util.concurrent.TimeoutException;
 import java.util.concurrent.atomic.AtomicReference;
 
 import static java.time.Duration.ofSeconds;
 import static java.util.Collections.singletonList;
 import static org.apache.commons.lang3.RandomStringUtils.randomAlphabetic;
-import static 
org.apache.flink.connector.pulsar.common.utils.PulsarExceptionUtils.sneakyAdmin;
 import static 
org.apache.flink.connector.pulsar.source.PulsarSourceOptions.PULSAR_ENABLE_AUTO_ACKNOWLEDGE_MESSAGE;
 import static 
org.apache.flink.connector.pulsar.source.PulsarSourceOptions.PULSAR_FETCH_ONE_MESSAGE_TIME;
 import static 
org.apache.flink.connector.pulsar.source.PulsarSourceOptions.PULSAR_MAX_FETCH_RECORDS;
@@ -67,7 +65,7 @@ import static org.assertj.core.api.Assertions.assertThat;
 class PulsarPartitionSplitReaderTest extends PulsarTestSuiteBase {
 
     @Test
-    void pollMessageAfterTimeout() throws InterruptedException, 
TimeoutException {
+    void pollMessageAfterTimeout() throws Exception {
         PulsarPartitionSplitReader splitReader = splitReader();
         String topicName = randomAlphabetic(10);
 
@@ -93,7 +91,7 @@ class PulsarPartitionSplitReaderTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void consumeMessageCreatedAfterHandleSplitChangesAndFetch() {
+    void consumeMessageCreatedAfterHandleSplitChangesAndFetch() throws 
Exception {
         PulsarPartitionSplitReader splitReader = splitReader();
         String topicName = randomAlphabetic(10);
 
@@ -103,7 +101,7 @@ class PulsarPartitionSplitReaderTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void consumeMessageCreatedBeforeHandleSplitsChanges() {
+    void consumeMessageCreatedBeforeHandleSplitsChanges() throws Exception {
         PulsarPartitionSplitReader splitReader = splitReader();
         String topicName = randomAlphabetic(10);
 
@@ -113,7 +111,8 @@ class PulsarPartitionSplitReaderTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void 
consumeMessageCreatedBeforeHandleSplitsChangesAndResetToEarliestPosition() {
+    void 
consumeMessageCreatedBeforeHandleSplitsChangesAndResetToEarliestPosition()
+            throws Exception {
         PulsarPartitionSplitReader splitReader = splitReader();
         String topicName = randomAlphabetic(10);
 
@@ -123,7 +122,7 @@ class PulsarPartitionSplitReaderTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void 
consumeMessageCreatedBeforeHandleSplitsChangesAndResetToLatestPosition() {
+    void 
consumeMessageCreatedBeforeHandleSplitsChangesAndResetToLatestPosition() throws 
Exception {
         PulsarPartitionSplitReader splitReader = splitReader();
         String topicName = randomAlphabetic(10);
 
@@ -133,20 +132,18 @@ class PulsarPartitionSplitReaderTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void 
consumeMessageCreatedBeforeHandleSplitsChangesAndUseSecondLastMessageIdCursor() 
{
+    void 
consumeMessageCreatedBeforeHandleSplitsChangesAndUseSecondLastMessageIdCursor()
+            throws Exception {
         PulsarPartitionSplitReader splitReader = splitReader();
         String topicName = randomAlphabetic(10);
 
         operator().setupTopic(topicName, STRING, () -> randomAlphabetic(10));
         MessageIdImpl lastMessageId =
                 (MessageIdImpl)
-                        sneakyAdmin(
-                                () ->
-                                        operator()
-                                                .admin()
-                                                .topics()
-                                                .getLastMessageId(
-                                                        
topicNameWithPartition(topicName, 0)));
+                        operator()
+                                .admin()
+                                .topics()
+                                
.getLastMessageId(topicNameWithPartition(topicName, 0));
         // when doing seek directly on consumer, by default it includes the 
specified messageId
         seekStartPositionAndHandleSplit(
                 splitReader,
@@ -160,7 +157,7 @@ class PulsarPartitionSplitReaderTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void emptyTopic() {
+    void emptyTopic() throws Exception {
         PulsarPartitionSplitReader splitReader = splitReader();
         String topicName = randomAlphabetic(10);
 
@@ -170,7 +167,7 @@ class PulsarPartitionSplitReaderTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void emptyTopicWithoutSeek() {
+    void emptyTopicWithoutSeek() throws Exception {
         PulsarPartitionSplitReader splitReader = splitReader();
         String topicName = randomAlphabetic(10);
 
@@ -211,7 +208,7 @@ class PulsarPartitionSplitReaderTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void consumeMessageCreatedBeforeHandleSplitsChangesWithoutSeek() {
+    void consumeMessageCreatedBeforeHandleSplitsChangesWithoutSeek() throws 
Exception {
         PulsarPartitionSplitReader splitReader = splitReader();
         String topicName = randomAlphabetic(10);
 
@@ -221,7 +218,8 @@ class PulsarPartitionSplitReaderTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void 
consumeMessageCreatedBeforeHandleSplitsChangesAndUseLatestStartCursorWithoutSeek()
 {
+    void 
consumeMessageCreatedBeforeHandleSplitsChangesAndUseLatestStartCursorWithoutSeek()
+            throws Exception {
         PulsarPartitionSplitReader splitReader = splitReader();
         String topicName = randomAlphabetic(10);
 
@@ -231,7 +229,8 @@ class PulsarPartitionSplitReaderTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void 
consumeMessageCreatedBeforeHandleSplitsChangesAndUseEarliestStartCursorWithoutSeek()
 {
+    void 
consumeMessageCreatedBeforeHandleSplitsChangesAndUseEarliestStartCursorWithoutSeek()
+            throws Exception {
         PulsarPartitionSplitReader splitReader = splitReader();
         String topicName = randomAlphabetic(10);
 
@@ -241,20 +240,18 @@ class PulsarPartitionSplitReaderTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void 
consumeMessageCreatedBeforeHandleSplitsChangesAndUseSecondLastMessageWithoutSeek()
 {
+    void 
consumeMessageCreatedBeforeHandleSplitsChangesAndUseSecondLastMessageWithoutSeek()
+            throws Exception {
         PulsarPartitionSplitReader splitReader = splitReader();
         String topicName = randomAlphabetic(10);
 
         operator().setupTopic(topicName, STRING, () -> randomAlphabetic(10));
         MessageIdImpl lastMessageId =
                 (MessageIdImpl)
-                        sneakyAdmin(
-                                () ->
-                                        operator()
-                                                .admin()
-                                                .topics()
-                                                .getLastMessageId(
-                                                        
topicNameWithPartition(topicName, 0)));
+                        operator()
+                                .admin()
+                                .topics()
+                                
.getLastMessageId(topicNameWithPartition(topicName, 0));
         // when recover, use exclusive startCursor
         handleSplit(
                 splitReader,
@@ -307,7 +304,7 @@ class PulsarPartitionSplitReaderTest extends 
PulsarTestSuiteBase {
     }
 
     private void seekStartPositionAndHandleSplit(
-            PulsarPartitionSplitReader reader, String topicName, int 
partitionId) {
+            PulsarPartitionSplitReader reader, String topicName, int 
partitionId) throws Exception {
         seekStartPositionAndHandleSplit(reader, topicName, partitionId, 
MessageId.latest);
     }
 
@@ -315,7 +312,8 @@ class PulsarPartitionSplitReaderTest extends 
PulsarTestSuiteBase {
             PulsarPartitionSplitReader reader,
             String topicName,
             int partitionId,
-            MessageId startPosition) {
+            MessageId startPosition)
+            throws Exception {
         TopicPartition partition = new TopicPartition(topicName, partitionId);
         PulsarPartitionSplit split =
                 new PulsarPartitionSplit(partition, StopCursor.never(), null, 
null);
@@ -326,23 +324,13 @@ class PulsarPartitionSplitReaderTest extends 
PulsarTestSuiteBase {
         SourceConfiguration sourceConfiguration = reader.sourceConfiguration;
         PulsarAdmin pulsarAdmin = reader.pulsarAdmin;
         String subscriptionName = sourceConfiguration.getSubscriptionName();
-        List<String> subscriptions =
-                sneakyAdmin(() -> 
pulsarAdmin.topics().getSubscriptions(topicName));
+        List<String> subscriptions = 
pulsarAdmin.topics().getSubscriptions(topicName);
         if (!subscriptions.contains(subscriptionName)) {
             // If this subscription is not available. Just create it.
-            sneakyAdmin(
-                    () ->
-                            pulsarAdmin
-                                    .topics()
-                                    .createSubscription(
-                                            topicName, subscriptionName, 
startPosition));
+            pulsarAdmin.topics().createSubscription(topicName, 
subscriptionName, startPosition);
         } else {
             // Reset the subscription if this is existed.
-            sneakyAdmin(
-                    () ->
-                            pulsarAdmin
-                                    .topics()
-                                    .resetCursor(topicName, subscriptionName, 
startPosition));
+            pulsarAdmin.topics().resetCursor(topicName, subscriptionName, 
startPosition);
         }
 
         // Accept the split and start consuming.
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/reader/PulsarSourceReaderTest.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/reader/PulsarSourceReaderTest.java
index a80f203..3ad03b8 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/reader/PulsarSourceReaderTest.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/reader/PulsarSourceReaderTest.java
@@ -218,7 +218,7 @@ class PulsarSourceReaderTest extends PulsarTestSuiteBase {
         reader.close();
     }
 
-    private String topicName() {
+    private String topicName() throws Exception {
         String topicName = randomAlphabetic(20);
         Random random = new Random(System.currentTimeMillis());
         operator().setupTopic(topicName, Schema.INT32, () -> 
random.nextInt(20));
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/reader/deserializer/PulsarDeserializationSchemaTest.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/reader/deserializer/PulsarDeserializationSchemaTest.java
index 32fe243..d20272d 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/reader/deserializer/PulsarDeserializationSchemaTest.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/source/reader/deserializer/PulsarDeserializationSchemaTest.java
@@ -127,7 +127,7 @@ class PulsarDeserializationSchemaTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void primitiveStringPulsarSchema() {
+    void primitiveStringPulsarSchema() throws Exception {
         final String topicName =
                 "primitiveString-" + ThreadLocalRandom.current().nextLong(0, 
Long.MAX_VALUE);
         operator().createTopic(topicName, 1);
@@ -143,7 +143,7 @@ class PulsarDeserializationSchemaTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void unversionedJsonStructPulsarSchema() {
+    void unversionedJsonStructPulsarSchema() throws Exception {
         final String topicName =
                 "unversionedJsonStruct-" + 
ThreadLocalRandom.current().nextLong(0, Long.MAX_VALUE);
         operator().createTopic(topicName, 1);
@@ -162,7 +162,7 @@ class PulsarDeserializationSchemaTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void keyValueJsonStructPulsarSchema() {
+    void keyValueJsonStructPulsarSchema() throws Exception {
         final String topicName =
                 "keyValueJsonStruct-" + 
ThreadLocalRandom.current().nextLong(0, Long.MAX_VALUE);
         operator().createTopic(topicName, 1);
@@ -187,7 +187,7 @@ class PulsarDeserializationSchemaTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void keyValueAvroStructPulsarSchema() {
+    void keyValueAvroStructPulsarSchema() throws Exception {
         final String topicName =
                 "keyValueAvroStruct-" + 
ThreadLocalRandom.current().nextLong(0, Long.MAX_VALUE);
         operator().createTopic(topicName, 1);
@@ -212,7 +212,7 @@ class PulsarDeserializationSchemaTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void keyValuePrimitivePulsarSchema() {
+    void keyValuePrimitivePulsarSchema() throws Exception {
         final String topicName =
                 "keyValuePrimitive-" + ThreadLocalRandom.current().nextLong(0, 
Long.MAX_VALUE);
         operator().createTopic(topicName, 1);
@@ -233,7 +233,7 @@ class PulsarDeserializationSchemaTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void keyValuePrimitiveKeyStructValuePulsarSchema() {
+    void keyValuePrimitiveKeyStructValuePulsarSchema() throws Exception {
         final String topicName =
                 "primitiveKeyStructValue-"
                         + ThreadLocalRandom.current().nextLong(0, 
Long.MAX_VALUE);
@@ -256,7 +256,7 @@ class PulsarDeserializationSchemaTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void keyValueStructKeyPrimitiveValuePulsarSchema() {
+    void keyValueStructKeyPrimitiveValuePulsarSchema() throws Exception {
         final String topicName =
                 "structKeyPrimitiveValue-"
                         + ThreadLocalRandom.current().nextLong(0, 
Long.MAX_VALUE);
@@ -279,7 +279,7 @@ class PulsarDeserializationSchemaTest extends 
PulsarTestSuiteBase {
     }
 
     @Test
-    void simpleFlinkSchema() {
+    void simpleFlinkSchema() throws Exception {
         final String topicName =
                 "simpleFlinkSchema-" + ThreadLocalRandom.current().nextLong(0, 
Long.MAX_VALUE);
         operator().createTopic(topicName, 1);
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/PulsarTestCommonUtils.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/PulsarTestCommonUtils.java
index 466fc3b..487fbdc 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/PulsarTestCommonUtils.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/PulsarTestCommonUtils.java
@@ -27,7 +27,6 @@ import org.apache.flink.streaming.api.CheckpointingMode;
 import org.apache.flink.test.resources.ResourceTestUtils;
 
 import org.apache.pulsar.client.api.MessageId;
-import org.junit.jupiter.api.extension.ParameterContext;
 
 import java.util.ArrayList;
 import java.util.List;
@@ -77,9 +76,4 @@ public class PulsarTestCommonUtils {
         }
         return splits;
     }
-
-    public static boolean isAssignableFromParameterContext(
-            Class<?> requiredType, ParameterContext context) {
-        return requiredType.isAssignableFrom(context.getParameter().getType());
-    }
 }
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/PulsarTestEnvironment.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/PulsarTestEnvironment.java
index f921e4b..33de6f4 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/PulsarTestEnvironment.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/PulsarTestEnvironment.java
@@ -86,25 +86,25 @@ public class PulsarTestEnvironment
 
     /** JUnit 5 Extension setup method. */
     @Override
-    public void beforeAll(ExtensionContext context) {
+    public void beforeAll(ExtensionContext context) throws Exception {
         runtime.startUp();
     }
 
     /** Start up the test resource. */
     @Override
-    public void startUp() {
+    public void startUp() throws Exception {
         runtime.startUp();
     }
 
     /** JUnit 5 Extension shutdown method. */
     @Override
-    public void afterAll(ExtensionContext context) {
+    public void afterAll(ExtensionContext context) throws Exception {
         runtime.tearDown();
     }
 
     /** Tear down the test resource. */
     @Override
-    public void tearDown() {
+    public void tearDown() throws Exception {
         runtime.tearDown();
     }
 
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/function/ControlSource.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/function/ControlSource.java
index 127e267..3ca23a9 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/function/ControlSource.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/function/ControlSource.java
@@ -32,7 +32,6 @@ import org.apache.flink.testutils.junit.SharedReference;
 import 
org.apache.flink.shaded.guava30.com.google.common.util.concurrent.Uninterruptibles;
 
 import org.apache.pulsar.client.api.Consumer;
-import org.apache.pulsar.client.api.ConsumerBuilder;
 import org.apache.pulsar.client.api.Message;
 import org.apache.pulsar.client.api.PulsarClientException;
 import org.apache.pulsar.client.api.Schema;
@@ -52,7 +51,6 @@ import java.util.concurrent.atomic.AtomicLong;
 import java.util.concurrent.atomic.AtomicReference;
 
 import static org.apache.commons.lang3.RandomStringUtils.randomAlphanumeric;
-import static 
org.apache.flink.connector.pulsar.common.utils.PulsarExceptionUtils.sneakyClient;
 import static org.apache.pulsar.client.api.SubscriptionMode.Durable;
 import static org.apache.pulsar.client.api.SubscriptionType.Exclusive;
 
@@ -77,7 +75,8 @@ public class ControlSource extends AbstractRichFunction
             DeliveryGuarantee guarantee,
             int messageCounts,
             Duration interval,
-            Duration timeout) {
+            Duration timeout)
+            throws PulsarClientException {
         MessageGenerator generator =
                 new MessageGenerator(topic, guarantee, messageCounts, 
interval);
         StopSignal signal = new StopSignal(operator, topic, messageCounts, 
timeout);
@@ -200,20 +199,21 @@ public class ControlSource extends AbstractRichFunction
         private final AtomicReference<PulsarClientException> 
throwableException;
 
         public StopSignal(
-                PulsarRuntimeOperator operator, String topic, int 
messageCounts, Duration timeout) {
+                PulsarRuntimeOperator operator, String topic, int 
messageCounts, Duration timeout)
+                throws PulsarClientException {
             this.desiredCounts = messageCounts;
             this.consumedRecords = Collections.synchronizedList(new 
ArrayList<>(messageCounts));
             this.deadline = new AtomicLong(timeout.toMillis() + 
System.currentTimeMillis());
             this.executor = Executors.newSingleThreadExecutor();
-            ConsumerBuilder<String> consumerBuilder =
+            this.consumer =
                     operator.client()
                             .newConsumer(Schema.STRING)
                             .topic(topic)
                             .subscriptionName(randomAlphanumeric(10))
                             .subscriptionMode(Durable)
                             .subscriptionType(Exclusive)
-                            
.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest);
-            this.consumer = sneakyClient(consumerBuilder::subscribe);
+                            
.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest)
+                            .subscribe();
             this.throwableException = new AtomicReference<>();
 
             // Start consuming.
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/runtime/PulsarRuntime.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/runtime/PulsarRuntime.java
index c8cd8e0..31a9e36 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/runtime/PulsarRuntime.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/runtime/PulsarRuntime.java
@@ -41,10 +41,10 @@ public interface PulsarRuntime {
     PulsarRuntime setConfigs(Map<String, String> configs);
 
     /** Start up this pulsar runtime, block the thread until everytime is 
ready for this runtime. */
-    void startUp();
+    void startUp() throws Exception;
 
     /** Shutdown this pulsar runtime. */
-    void tearDown();
+    void tearDown() throws Exception;
 
     /**
      * Return an operator for operating this pulsar runtime. This operator 
predefined a set of
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/runtime/PulsarRuntimeOperator.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/runtime/PulsarRuntimeOperator.java
index 310462f..6fe4d29 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/runtime/PulsarRuntimeOperator.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/runtime/PulsarRuntimeOperator.java
@@ -30,13 +30,10 @@ import org.apache.pulsar.client.admin.PulsarAdminException;
 import org.apache.pulsar.client.admin.PulsarAdminException.ConflictException;
 import org.apache.pulsar.client.admin.PulsarAdminException.NotFoundException;
 import org.apache.pulsar.client.api.Consumer;
-import org.apache.pulsar.client.api.ConsumerBuilder;
 import org.apache.pulsar.client.api.Message;
 import org.apache.pulsar.client.api.MessageId;
 import org.apache.pulsar.client.api.Producer;
-import org.apache.pulsar.client.api.ProducerBuilder;
 import org.apache.pulsar.client.api.PulsarClient;
-import org.apache.pulsar.client.api.PulsarClientException;
 import org.apache.pulsar.client.api.Schema;
 import org.apache.pulsar.client.api.TypedMessageBuilder;
 import org.apache.pulsar.client.api.transaction.TransactionCoordinatorClient;
@@ -45,12 +42,12 @@ import org.apache.pulsar.common.naming.TopicName;
 import org.apache.pulsar.common.partition.PartitionedTopicMetadata;
 
 import java.io.Closeable;
+import java.io.IOException;
 import java.time.Duration;
 import java.util.ArrayList;
 import java.util.Collection;
 import java.util.List;
 import java.util.Random;
-import java.util.concurrent.ExecutionException;
 import java.util.function.Supplier;
 import java.util.stream.Stream;
 
@@ -63,9 +60,6 @@ import static 
org.apache.flink.connector.base.DeliveryGuarantee.EXACTLY_ONCE;
 import static 
org.apache.flink.connector.pulsar.common.config.PulsarOptions.PULSAR_ADMIN_URL;
 import static 
org.apache.flink.connector.pulsar.common.config.PulsarOptions.PULSAR_ENABLE_TRANSACTION;
 import static 
org.apache.flink.connector.pulsar.common.config.PulsarOptions.PULSAR_SERVICE_URL;
-import static 
org.apache.flink.connector.pulsar.common.utils.PulsarExceptionUtils.sneakyAdmin;
-import static 
org.apache.flink.connector.pulsar.common.utils.PulsarExceptionUtils.sneakyClient;
-import static 
org.apache.flink.connector.pulsar.common.utils.PulsarExceptionUtils.sneakyThrow;
 import static 
org.apache.flink.connector.pulsar.common.utils.PulsarTransactionUtils.getTcClient;
 import static 
org.apache.flink.connector.pulsar.sink.PulsarSinkOptions.PULSAR_SEND_TIMEOUT_MS;
 import static 
org.apache.flink.connector.pulsar.sink.PulsarSinkOptions.PULSAR_WRITE_DELIVERY_GUARANTEE;
@@ -91,7 +85,7 @@ public class PulsarRuntimeOperator implements Closeable {
     private final PulsarClient client;
     private final PulsarAdmin admin;
 
-    public PulsarRuntimeOperator(String serviceUrl, String adminUrl) {
+    public PulsarRuntimeOperator(String serviceUrl, String adminUrl) throws 
Exception {
         this(serviceUrl, serviceUrl, adminUrl, adminUrl);
     }
 
@@ -99,17 +93,12 @@ public class PulsarRuntimeOperator implements Closeable {
             String serviceUrl,
             String containerServiceUrl,
             String adminUrl,
-            String containerAdminUrl) {
+            String containerAdminUrl)
+            throws Exception {
         this.serviceUrl = containerServiceUrl;
         this.adminUrl = containerAdminUrl;
-        this.client =
-                sneakyClient(
-                        () ->
-                                PulsarClient.builder()
-                                        .serviceUrl(serviceUrl)
-                                        .enableTransaction(true)
-                                        .build());
-        this.admin = sneakyClient(() -> 
PulsarAdmin.builder().serviceHttpUrl(adminUrl).build());
+        this.client = 
PulsarClient.builder().serviceUrl(serviceUrl).enableTransaction(true).build();
+        this.admin = PulsarAdmin.builder().serviceHttpUrl(adminUrl).build();
     }
 
     /**
@@ -118,7 +107,7 @@ public class PulsarRuntimeOperator implements Closeable {
      *
      * @param topic Pulsar topic name, it couldn't be a name with partition 
index.
      */
-    public void setupTopic(String topic) {
+    public void setupTopic(String topic) throws Exception {
         Random random = new Random(System.currentTimeMillis());
         setupTopic(topic, Schema.STRING, () -> randomAlphanumeric(10 + 
random.nextInt(20)));
     }
@@ -131,7 +120,8 @@ public class PulsarRuntimeOperator implements Closeable {
      * @param schema The Pulsar schema for serializing records into bytes.
      * @param supplier The supplier for providing the records which would be 
sent to Pulsar.
      */
-    public <T> void setupTopic(String topic, Schema<T> schema, Supplier<T> 
supplier) {
+    public <T> void setupTopic(String topic, Schema<T> schema, Supplier<T> 
supplier)
+            throws Exception {
         setupTopic(topic, schema, supplier, NUM_RECORDS_PER_PARTITION);
     }
 
@@ -145,7 +135,8 @@ public class PulsarRuntimeOperator implements Closeable {
      * @param numRecordsPerSplit The number of records for a partition.
      */
     public <T> void setupTopic(
-            String topic, Schema<T> schema, Supplier<T> supplier, int 
numRecordsPerSplit) {
+            String topic, Schema<T> schema, Supplier<T> supplier, int 
numRecordsPerSplit)
+            throws Exception {
         String topicName = topicName(topic);
         createTopic(topicName, DEFAULT_PARTITIONS);
 
@@ -167,7 +158,7 @@ public class PulsarRuntimeOperator implements Closeable {
      * @param numberOfPartitions The number of partitions. We would create a 
non-partitioned topic
      *     if this number is zero.
      */
-    public void createTopic(String topic, int numberOfPartitions) {
+    public void createTopic(String topic, int numberOfPartitions) throws 
Exception {
         checkArgument(numberOfPartitions >= 0);
         if (numberOfPartitions == 0) {
             createNonPartitionedTopic(topic);
@@ -176,8 +167,8 @@ public class PulsarRuntimeOperator implements Closeable {
         }
     }
 
-    public void createSchema(String topic, Schema<?> schema) {
-        sneakyAdmin(() -> admin().schemas().createSchema(topic, 
schema.getSchemaInfo()));
+    public void createSchema(String topic, Schema<?> schema) throws Exception {
+        admin().schemas().createSchema(topic, schema.getSchemaInfo());
     }
 
     /**
@@ -186,14 +177,13 @@ public class PulsarRuntimeOperator implements Closeable {
      * @param topic The topic name.
      * @param newPartitionsNum The new partition size which should exceed 
previous size.
      */
-    public void increaseTopicPartitions(String topic, int newPartitionsNum) {
-        PartitionedTopicMetadata metadata =
-                sneakyAdmin(() -> 
admin().topics().getPartitionedTopicMetadata(topic));
+    public void increaseTopicPartitions(String topic, int newPartitionsNum) 
throws Exception {
+        PartitionedTopicMetadata metadata = 
admin().topics().getPartitionedTopicMetadata(topic);
         checkArgument(
                 metadata.partitions < newPartitionsNum,
                 "The new partition size which should greater than previous 
size.");
 
-        sneakyAdmin(() -> admin().topics().updatePartitionedTopic(topic, 
newPartitionsNum));
+        admin().topics().updatePartitionedTopic(topic, newPartitionsNum);
     }
 
     /**
@@ -201,7 +191,7 @@ public class PulsarRuntimeOperator implements Closeable {
      *
      * @param topic The topic name.
      */
-    public void deleteTopic(String topic) {
+    public void deleteTopic(String topic) throws Exception {
         String topicName = topicName(topic);
         PartitionedTopicMetadata metadata;
 
@@ -210,27 +200,20 @@ public class PulsarRuntimeOperator implements Closeable {
         } catch (NotFoundException e) {
             // This topic doesn't exist. Just skip deletion.
             return;
-        } catch (PulsarAdminException e) {
-            sneakyThrow(e);
-            return;
         }
 
         if (metadata.partitions == NON_PARTITIONED) {
-            sneakyAdmin(() -> admin().topics().delete(topicName));
+            admin().topics().delete(topicName);
         } else {
-            sneakyAdmin(() -> 
admin().topics().deletePartitionedTopic(topicName));
+            admin().topics().deletePartitionedTopic(topicName);
         }
     }
 
     /** Convert the topic metadata into a list of topic partitions. */
-    public List<TopicPartition> topicInfo(String topic) {
-        try {
-            return client().getPartitionsForTopic(topic).get().stream()
-                    .map(p -> new TopicPartition(topic, 
TopicName.getPartitionIndex(p)))
-                    .collect(toList());
-        } catch (InterruptedException | ExecutionException e) {
-            throw new IllegalStateException(e);
-        }
+    public List<TopicPartition> topicInfo(String topic) throws Exception {
+        return client().getPartitionsForTopic(topic).get().stream()
+                .map(p -> new TopicPartition(topic, 
TopicName.getPartitionIndex(p)))
+                .collect(toList());
     }
 
     /**
@@ -242,7 +225,7 @@ public class PulsarRuntimeOperator implements Closeable {
      * @param <T> The type of the record.
      * @return message id.
      */
-    public <T> MessageId sendMessage(String topic, Schema<T> schema, T 
message) {
+    public <T> MessageId sendMessage(String topic, Schema<T> schema, T 
message) throws Exception {
         List<MessageId> messageIds = sendMessages(topic, schema, 
singletonList(message));
         checkArgument(messageIds.size() == 1);
 
@@ -259,7 +242,8 @@ public class PulsarRuntimeOperator implements Closeable {
      * @param <T> The type of the record.
      * @return message id.
      */
-    public <T> MessageId sendMessage(String topic, Schema<T> schema, String 
key, T message) {
+    public <T> MessageId sendMessage(String topic, Schema<T> schema, String 
key, T message)
+            throws Exception {
         List<MessageId> messageIds = sendMessages(topic, schema, key, 
singletonList(message));
         checkArgument(messageIds.size() == 1);
 
@@ -275,8 +259,8 @@ public class PulsarRuntimeOperator implements Closeable {
      * @param <T> The type of the record.
      * @return message id.
      */
-    public <T> List<MessageId> sendMessages(
-            String topic, Schema<T> schema, Collection<T> messages) {
+    public <T> List<MessageId> sendMessages(String topic, Schema<T> schema, 
Collection<T> messages)
+            throws Exception {
         return sendMessages(topic, schema, null, messages);
     }
 
@@ -291,10 +275,9 @@ public class PulsarRuntimeOperator implements Closeable {
      * @return message id.
      */
     public <T> List<MessageId> sendMessages(
-            String topic, Schema<T> schema, String key, Collection<T> 
messages) {
+            String topic, Schema<T> schema, String key, Collection<T> 
messages) throws Exception {
         try (Producer<T> producer = createProducer(topic, schema)) {
             List<MessageId> messageIds = new ArrayList<>(messages.size());
-
             for (T message : messages) {
                 TypedMessageBuilder<T> builder = 
producer.newMessage().value(message);
                 if (!Strings.isNullOrEmpty(key)) {
@@ -305,9 +288,6 @@ public class PulsarRuntimeOperator implements Closeable {
             }
             producer.flush();
             return messageIds;
-        } catch (PulsarClientException e) {
-            sneakyThrow(e);
-            return emptyList();
         }
     }
 
@@ -315,14 +295,11 @@ public class PulsarRuntimeOperator implements Closeable {
      * Consume a message from the given Pulsar topic, this method would be 
blocked until we get a
      * message from this topic.
      */
-    public <T> Message<T> receiveMessage(String topic, Schema<T> schema) {
+    public <T> Message<T> receiveMessage(String topic, Schema<T> schema) 
throws Exception {
         try (Consumer<T> consumer = createConsumer(topic, schema)) {
             Message<T> message = consumer.receive();
             consumer.acknowledge(message.getMessageId());
             return message;
-        } catch (PulsarClientException e) {
-            sneakyThrow(e);
-            return null;
         }
     }
 
@@ -335,7 +312,6 @@ public class PulsarRuntimeOperator implements Closeable {
             Message<T> message =
                     consumer.receive(Math.toIntExact(timeout.toMillis()), 
MILLISECONDS);
             consumer.acknowledge(message.getMessageId());
-
             return message;
         } catch (Exception e) {
             return null;
@@ -346,7 +322,8 @@ public class PulsarRuntimeOperator implements Closeable {
      * Consume a fixed number of messages from the given Pulsar topic, this 
method would be blocked
      * until we get the exactly number of messages from this topic.
      */
-    public <T> List<Message<T>> receiveMessages(String topic, Schema<T> 
schema, int counts) {
+    public <T> List<Message<T>> receiveMessages(String topic, Schema<T> 
schema, int counts)
+            throws Exception {
         if (counts == 0) {
             return emptyList();
         } else if (counts < 0) {
@@ -365,10 +342,8 @@ public class PulsarRuntimeOperator implements Closeable {
                     messages.add(message);
                     consumer.acknowledge(message.getMessageId());
                 }
+
                 return messages;
-            } catch (PulsarClientException e) {
-                sneakyThrow(e);
-                return emptyList();
             }
         }
     }
@@ -448,7 +423,7 @@ public class PulsarRuntimeOperator implements Closeable {
      * manually.
      */
     @Override
-    public void close() throws PulsarClientException {
+    public void close() throws IOException {
         if (admin != null) {
             admin.close();
         }
@@ -459,55 +434,50 @@ public class PulsarRuntimeOperator implements Closeable {
 
     // --------------------------- Private Methods 
-----------------------------
 
-    private void createNonPartitionedTopic(String topic) {
+    private void createNonPartitionedTopic(String topic) throws Exception {
         try {
             admin().topics().createNonPartitionedTopic(topic);
         } catch (PulsarAdminException e) {
             if (!(e instanceof ConflictException
                     && e.getMessage().equals("This topic already exists"))) {
-                sneakyThrow(e);
+                throw e;
             }
         }
     }
 
-    private void createPartitionedTopic(String topic, int numberOfPartitions) {
+    private void createPartitionedTopic(String topic, int numberOfPartitions) 
throws Exception {
         try {
             admin().topics().createPartitionedTopic(topic, numberOfPartitions);
         } catch (PulsarAdminException e) {
             if (!(e instanceof ConflictException
                     && e.getMessage().equals("This topic already exists"))) {
-                sneakyThrow(e);
+                throw e;
             }
         }
     }
 
-    private <T> Producer<T> createProducer(String topic, Schema<T> schema) {
-        ProducerBuilder<T> builder =
-                client().newProducer(schema)
-                        .topic(topic)
-                        .enableBatching(false)
-                        .enableMultiSchema(true)
-                        .accessMode(Shared);
-
-        return sneakyClient(builder::create);
+    private <T> Producer<T> createProducer(String topic, Schema<T> schema) 
throws Exception {
+        return client().newProducer(schema)
+                .topic(topic)
+                .enableBatching(false)
+                .enableMultiSchema(true)
+                .accessMode(Shared)
+                .create();
     }
 
-    private <T> Consumer<T> createConsumer(String topic, Schema<T> schema) {
+    private <T> Consumer<T> createConsumer(String topic, Schema<T> schema) 
throws Exception {
         // Create the earliest subscription if it's not existed.
-        List<String> subscriptions = sneakyAdmin(() -> 
admin().topics().getSubscriptions(topic));
+        List<String> subscriptions = admin().topics().getSubscriptions(topic);
         if (!subscriptions.contains(SUBSCRIPTION_NAME)) {
-            sneakyAdmin(
-                    () -> admin().topics().createSubscription(topic, 
SUBSCRIPTION_NAME, earliest));
+            admin().topics().createSubscription(topic, SUBSCRIPTION_NAME, 
earliest);
         }
 
         // Create the consumer without the initial position.
-        ConsumerBuilder<T> builder =
-                client().newConsumer(schema)
-                        .topic(topic)
-                        .subscriptionName(SUBSCRIPTION_NAME)
-                        .subscriptionMode(Durable)
-                        .subscriptionType(Exclusive);
-
-        return sneakyClient(builder::subscribe);
+        return client().newConsumer(schema)
+                .topic(topic)
+                .subscriptionName(SUBSCRIPTION_NAME)
+                .subscriptionMode(Durable)
+                .subscriptionType(Exclusive)
+                .subscribe();
     }
 }
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/runtime/container/PulsarContainerRuntime.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/runtime/container/PulsarContainerRuntime.java
index c709432..c77aa1e 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/runtime/container/PulsarContainerRuntime.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/runtime/container/PulsarContainerRuntime.java
@@ -101,7 +101,7 @@ public class PulsarContainerRuntime implements 
PulsarRuntime {
     }
 
     @Override
-    public void startUp() {
+    public void startUp() throws Exception {
         if (!started.compareAndSet(false, true)) {
             LOG.warn("You have started the Pulsar Container. We will skip this 
execution.");
             return;
@@ -144,16 +144,12 @@ public class PulsarContainerRuntime implements 
PulsarRuntime {
     }
 
     @Override
-    public void tearDown() {
-        try {
-            if (operator != null) {
-                operator.close();
-            }
-            container.stop();
-            started.compareAndSet(true, false);
-        } catch (Exception e) {
-            throw new IllegalStateException(e);
+    public void tearDown() throws Exception {
+        if (operator != null) {
+            operator.close();
         }
+        container.stop();
+        started.compareAndSet(true, false);
     }
 
     @Override
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/sink/PulsarSinkTestContext.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/sink/PulsarSinkTestContext.java
index 8c6877c..3bc7c34 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/sink/PulsarSinkTestContext.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/sink/PulsarSinkTestContext.java
@@ -15,6 +15,7 @@ import 
org.apache.flink.connector.pulsar.testutils.sink.reader.PulsarPartitionDa
 import org.apache.flink.connector.testframe.external.ExternalSystemDataReader;
 import 
org.apache.flink.connector.testframe.external.sink.DataStreamSinkV2ExternalContext;
 import org.apache.flink.connector.testframe.external.sink.TestingSinkSettings;
+import org.apache.flink.util.FlinkRuntimeException;
 
 import org.apache.flink.shaded.guava30.com.google.common.io.Closer;
 
@@ -55,7 +56,11 @@ public abstract class PulsarSinkTestContext extends 
PulsarTestContext<String>
         // Create the topic if it needs.
         if (creatTopic()) {
             for (String topic : topics) {
-                operator.createTopic(topic, 4);
+                try {
+                    operator.createTopic(topic, 4);
+                } catch (Exception e) {
+                    throw new FlinkRuntimeException(e);
+                }
             }
         }
 
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/sink/reader/PulsarPartitionDataReader.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/sink/reader/PulsarPartitionDataReader.java
index 55c37b0..ae2af9d 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/sink/reader/PulsarPartitionDataReader.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/sink/reader/PulsarPartitionDataReader.java
@@ -21,6 +21,7 @@ package 
org.apache.flink.connector.pulsar.testutils.sink.reader;
 import org.apache.flink.connector.pulsar.common.crypto.PulsarCrypto;
 import 
org.apache.flink.connector.pulsar.testutils.runtime.PulsarRuntimeOperator;
 import org.apache.flink.connector.testframe.external.ExternalSystemDataReader;
+import org.apache.flink.util.FlinkRuntimeException;
 
 import org.apache.pulsar.client.api.Consumer;
 import org.apache.pulsar.client.api.ConsumerBuilder;
@@ -44,7 +45,6 @@ import java.util.List;
 
 import static java.util.concurrent.TimeUnit.MILLISECONDS;
 import static org.apache.commons.lang3.RandomStringUtils.randomAlphanumeric;
-import static 
org.apache.flink.connector.pulsar.common.utils.PulsarExceptionUtils.sneakyClient;
 
 /** The data reader for a specified topic partition from Pulsar. */
 public class PulsarPartitionDataReader<T> implements 
ExternalSystemDataReader<T>, Closeable {
@@ -52,7 +52,6 @@ public class PulsarPartitionDataReader<T> implements 
ExternalSystemDataReader<T>
     private static final Logger LOG = 
LoggerFactory.getLogger(PulsarPartitionDataReader.class);
 
     private final Consumer<T> consumer;
-    private final PulsarRuntimeOperator operator;
 
     public PulsarPartitionDataReader(
             PulsarRuntimeOperator operator, List<String> topics, Schema<T> 
schema) {
@@ -87,8 +86,11 @@ public class PulsarPartitionDataReader<T> implements 
ExternalSystemDataReader<T>
             }
         }
 
-        this.consumer = sneakyClient(builder::subscribe);
-        this.operator = operator;
+        try {
+            this.consumer = builder.subscribe();
+        } catch (PulsarClientException e) {
+            throw new FlinkRuntimeException(e);
+        }
     }
 
     @Override
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/cases/MultipleTopicsConsumingContext.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/cases/MultipleTopicsConsumingContext.java
index 009a98c..a62a856 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/cases/MultipleTopicsConsumingContext.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/cases/MultipleTopicsConsumingContext.java
@@ -20,6 +20,7 @@ package 
org.apache.flink.connector.pulsar.testutils.source.cases;
 
 import org.apache.flink.connector.pulsar.testutils.PulsarTestEnvironment;
 import 
org.apache.flink.connector.pulsar.testutils.source.PulsarSourceTestContext;
+import org.apache.flink.util.FlinkRuntimeException;
 
 import static org.apache.commons.lang3.RandomStringUtils.randomAlphabetic;
 import static 
org.apache.flink.connector.pulsar.source.enumerator.topic.TopicNameUtils.topicNameWithPartition;
@@ -30,7 +31,8 @@ import static 
org.apache.flink.connector.pulsar.source.enumerator.topic.TopicNam
  */
 public class MultipleTopicsConsumingContext extends PulsarSourceTestContext {
 
-    private final String topicPrefix = "flink-multiple-topic-" + 
randomAlphabetic(8) + "-";
+    private final String topicPrefix =
+            "public/default/flink-multiple-topic-" + randomAlphabetic(8) + "-";
 
     private int index = 0;
 
@@ -45,7 +47,7 @@ public class MultipleTopicsConsumingContext extends 
PulsarSourceTestContext {
 
     @Override
     protected String topicPattern() {
-        return topicPrefix + ".+";
+        return topicPrefix + "\\d+";
     }
 
     @Override
@@ -56,7 +58,11 @@ public class MultipleTopicsConsumingContext extends 
PulsarSourceTestContext {
     @Override
     protected String generatePartitionName() {
         String topic = topicPrefix + index;
-        operator.createTopic(topic, 1);
+        try {
+            operator.createTopic(topic, 1);
+        } catch (Exception e) {
+            throw new FlinkRuntimeException(e);
+        }
         index++;
 
         return topicNameWithPartition(topic, 0);
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/cases/SingleTopicConsumingContext.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/cases/SingleTopicConsumingContext.java
index 9755e58..dd850fc 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/cases/SingleTopicConsumingContext.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/cases/SingleTopicConsumingContext.java
@@ -20,6 +20,7 @@ package 
org.apache.flink.connector.pulsar.testutils.source.cases;
 
 import org.apache.flink.connector.pulsar.testutils.PulsarTestEnvironment;
 import 
org.apache.flink.connector.pulsar.testutils.source.PulsarSourceTestContext;
+import org.apache.flink.util.FlinkRuntimeException;
 
 import static org.apache.commons.lang3.RandomStringUtils.randomAlphanumeric;
 import static 
org.apache.flink.connector.pulsar.source.enumerator.topic.TopicNameUtils.topicNameWithPartition;
@@ -30,7 +31,7 @@ import static 
org.apache.flink.connector.pulsar.source.enumerator.topic.TopicNam
  */
 public class SingleTopicConsumingContext extends PulsarSourceTestContext {
 
-    private final String topicName = "pulsar-single-topic-" + 
randomAlphanumeric(8);
+    private final String topicName = "public/default/pulsar-single-topic-" + 
randomAlphanumeric(8);
 
     private int index = 0;
 
@@ -45,7 +46,7 @@ public class SingleTopicConsumingContext extends 
PulsarSourceTestContext {
 
     @Override
     protected String topicPattern() {
-        return topicName + ".+";
+        return topicName;
     }
 
     @Override
@@ -55,10 +56,14 @@ public class SingleTopicConsumingContext extends 
PulsarSourceTestContext {
 
     @Override
     protected String generatePartitionName() {
-        if (index == 0) {
-            operator.createTopic(topicName, index + 1);
-        } else {
-            operator.increaseTopicPartitions(topicName, index + 1);
+        try {
+            if (index == 0) {
+                operator.createTopic(topicName, index + 1);
+            } else {
+                operator.increaseTopicPartitions(topicName, index + 1);
+            }
+        } catch (Exception e) {
+            throw new FlinkRuntimeException(e);
         }
 
         return topicNameWithPartition(topicName, index++);
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/writer/KeyedPulsarPartitionDataWriter.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/writer/KeyedPulsarPartitionDataWriter.java
index 194aa48..6d3c9af 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/writer/KeyedPulsarPartitionDataWriter.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/writer/KeyedPulsarPartitionDataWriter.java
@@ -20,6 +20,7 @@ package 
org.apache.flink.connector.pulsar.testutils.source.writer;
 
 import 
org.apache.flink.connector.pulsar.testutils.runtime.PulsarRuntimeOperator;
 import 
org.apache.flink.connector.testframe.external.ExternalSystemSplitDataWriter;
+import org.apache.flink.util.FlinkRuntimeException;
 
 import org.apache.pulsar.client.api.Schema;
 
@@ -51,12 +52,16 @@ public class KeyedPulsarPartitionDataWriter implements 
ExternalSystemSplitDataWr
 
     @Override
     public void writeRecords(List<String> records) {
-        // Send messages with the key we don't need.
-        List<String> newRecords = records.stream().map(a -> a + 
keyToRead).collect(toList());
-        operator.sendMessages(fullTopicName, Schema.STRING, keyToExclude, 
newRecords);
+        try {
+            // Send messages with the key we don't need.
+            List<String> newRecords = records.stream().map(a -> a + 
keyToRead).collect(toList());
+            operator.sendMessages(fullTopicName, Schema.STRING, keyToExclude, 
newRecords);
 
-        // Send messages with the given key.
-        operator.sendMessages(fullTopicName, Schema.STRING, keyToRead, 
records);
+            // Send messages with the given key.
+            operator.sendMessages(fullTopicName, Schema.STRING, keyToRead, 
records);
+        } catch (Exception e) {
+            throw new FlinkRuntimeException(e);
+        }
     }
 
     @Override
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/writer/PulsarEncryptDataWriter.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/writer/PulsarEncryptDataWriter.java
index 9eeb8c9..1743c7e 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/writer/PulsarEncryptDataWriter.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/writer/PulsarEncryptDataWriter.java
@@ -21,11 +21,13 @@ package 
org.apache.flink.connector.pulsar.testutils.source.writer;
 import org.apache.flink.connector.pulsar.common.crypto.PulsarCrypto;
 import 
org.apache.flink.connector.pulsar.testutils.runtime.PulsarRuntimeOperator;
 import 
org.apache.flink.connector.testframe.external.ExternalSystemSplitDataWriter;
+import org.apache.flink.util.FlinkRuntimeException;
 
 import org.apache.pulsar.client.api.MessageCrypto;
 import org.apache.pulsar.client.api.Producer;
 import org.apache.pulsar.client.api.ProducerBuilder;
 import org.apache.pulsar.client.api.ProducerCryptoFailureAction;
+import org.apache.pulsar.client.api.PulsarClientException;
 import org.apache.pulsar.client.api.Schema;
 import org.apache.pulsar.client.impl.ProducerBuilderImpl;
 import org.apache.pulsar.client.impl.conf.ProducerConfigurationData;
@@ -34,7 +36,6 @@ import org.apache.pulsar.common.api.proto.MessageMetadata;
 import java.util.List;
 import java.util.Set;
 
-import static 
org.apache.flink.connector.pulsar.common.utils.PulsarExceptionUtils.sneakyClient;
 import static org.apache.pulsar.client.api.ProducerAccessMode.Shared;
 
 /** Encrypt the messages with the given public key and send the message to 
Pulsar. */
@@ -67,15 +68,23 @@ public class PulsarEncryptDataWriter<T> implements 
ExternalSystemSplitDataWriter
             conf.setMessageCrypto(messageCrypto);
         }
 
-        this.producer = sneakyClient(builder::create);
+        try {
+            this.producer = builder.create();
+        } catch (PulsarClientException e) {
+            throw new FlinkRuntimeException(e);
+        }
     }
 
     @Override
     public void writeRecords(List<T> records) {
-        for (T record : records) {
-            sneakyClient(() -> producer.newMessage().value(record).send());
+        try {
+            for (T record : records) {
+                producer.newMessage().value(record).send();
+            }
+            producer.flush();
+        } catch (PulsarClientException e) {
+            throw new FlinkRuntimeException(e);
         }
-        sneakyClient(producer::flush);
     }
 
     @Override
diff --git 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/writer/PulsarPartitionDataWriter.java
 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/writer/PulsarPartitionDataWriter.java
index c5893c9..2d53019 100644
--- 
a/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/writer/PulsarPartitionDataWriter.java
+++ 
b/flink-connector-pulsar/src/test/java/org/apache/flink/connector/pulsar/testutils/source/writer/PulsarPartitionDataWriter.java
@@ -20,6 +20,7 @@ package 
org.apache.flink.connector.pulsar.testutils.source.writer;
 
 import 
org.apache.flink.connector.pulsar.testutils.runtime.PulsarRuntimeOperator;
 import 
org.apache.flink.connector.testframe.external.ExternalSystemSplitDataWriter;
+import org.apache.flink.util.FlinkRuntimeException;
 
 import org.apache.pulsar.client.api.Schema;
 
@@ -44,7 +45,11 @@ public class PulsarPartitionDataWriter<T> implements 
ExternalSystemSplitDataWrit
 
     @Override
     public void writeRecords(List<T> records) {
-        operator.sendMessages(fullTopicName, schema, records);
+        try {
+            operator.sendMessages(fullTopicName, schema, records);
+        } catch (Exception e) {
+            throw new FlinkRuntimeException(e);
+        }
     }
 
     @Override

Reply via email to