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