From 1ec2641bfd237845fc2b57fba2b38cf699ae6964 Mon Sep 17 00:00:00 2001 From: Ming Luo Date: Wed, 27 May 2020 23:38:34 -0400 Subject: [PATCH] auth token for debezium and kafka connect adaptor --- .../io/debezium/PulsarDatabaseHistory.java | 22 ++++++++++++++++--- .../connect/PulsarKafkaWorkerConfig.java | 10 +++++++++ .../connect/PulsarOffsetBackingStore.java | 14 +++++++++--- 3 files changed, 40 insertions(+), 6 deletions(-) diff --git a/pulsar-io/debezium/core/src/main/java/org/apache/pulsar/io/debezium/PulsarDatabaseHistory.java b/pulsar-io/debezium/core/src/main/java/org/apache/pulsar/io/debezium/PulsarDatabaseHistory.java index c3d95a7a6dedf..46878c150a19a 100644 --- a/pulsar-io/debezium/core/src/main/java/org/apache/pulsar/io/debezium/PulsarDatabaseHistory.java +++ b/pulsar-io/debezium/core/src/main/java/org/apache/pulsar/io/debezium/PulsarDatabaseHistory.java @@ -36,6 +36,8 @@ import org.apache.kafka.common.config.ConfigDef.Importance; import org.apache.kafka.common.config.ConfigDef.Type; import org.apache.kafka.common.config.ConfigDef.Width; +import org.apache.pulsar.client.api.AuthenticationFactory; +import org.apache.pulsar.client.api.ClientBuilder; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Producer; @@ -68,14 +70,24 @@ public final class PulsarDatabaseHistory extends AbstractDatabaseHistory { .withDescription("Pulsar service url") .withValidation(Field::isRequired); + public static final Field PULSAR_TOKEN = Field.create(CONFIGURATION_FIELD_PREFIX_STRING + "pulsar.token") + .withDisplayName("Pulsar auth token") + .withType(Type.STRING) + .withWidth(Width.LONG) + .withImportance(Importance.HIGH) + .withDescription("Pulsar authentication token") + .withValidation(Field::isOptional); + public static Field.Set ALL_FIELDS = Field.setOf( TOPIC, SERVICE_URL, + PULSAR_TOKEN, DatabaseHistory.NAME); private final DocumentReader reader = DocumentReader.defaultReader(); private String topicName; private String serviceUrl; + private String token; private String dbHistoryName; private volatile PulsarClient pulsarClient; private volatile Producer producer; @@ -94,6 +106,7 @@ public void configure( } this.topicName = config.getString(TOPIC); this.serviceUrl = config.getString(SERVICE_URL); + this.token = config.getString(PULSAR_TOKEN); // Copy the relevant portions of the configuration and add useful defaults ... this.dbHistoryName = config.getString(DatabaseHistory.NAME, UUID.randomUUID().toString()); @@ -117,9 +130,12 @@ public void initializeStorage() { void setupClientIfNeeded() { if (null == this.pulsarClient) { try { - pulsarClient = PulsarClient.builder() - .serviceUrl(serviceUrl) - .build(); + ClientBuilder builder = PulsarClient.builder().serviceUrl(serviceUrl); + + if (token != null && token != "") { + builder = builder.authentication(AuthenticationFactory.token(token)); + } + pulsarClient = builder.build(); } catch (PulsarClientException e) { throw new RuntimeException("Failed to create pulsar client to pulsar cluster at " + serviceUrl, e); diff --git a/pulsar-io/kafka-connect-adaptor/src/main/java/org/apache/pulsar/io/kafka/connect/PulsarKafkaWorkerConfig.java b/pulsar-io/kafka-connect-adaptor/src/main/java/org/apache/pulsar/io/kafka/connect/PulsarKafkaWorkerConfig.java index 624c59acd84de..920927413ba4c 100644 --- a/pulsar-io/kafka-connect-adaptor/src/main/java/org/apache/pulsar/io/kafka/connect/PulsarKafkaWorkerConfig.java +++ b/pulsar-io/kafka-connect-adaptor/src/main/java/org/apache/pulsar/io/kafka/connect/PulsarKafkaWorkerConfig.java @@ -44,6 +44,12 @@ public class PulsarKafkaWorkerConfig extends WorkerConfig { public static final String PULSAR_SERVICE_URL_CONFIG = "pulsar.service.url"; private static final String PULSAR_SERVICE_URL_CONFIG_DOC = "pulsar service url"; + /** + * pulsar.auth.token + */ + public static final String PULSAR_AUTH_TOKEN_CONFIG = "pulsar.auth.token"; + private static final String PULSAR_AUTH_TOKEN_CONFIG_DOC = "pulsar auth token"; + /** * topic.namespace */ @@ -60,6 +66,10 @@ public class PulsarKafkaWorkerConfig extends WorkerConfig { Type.STRING, Importance.HIGH, PULSAR_SERVICE_URL_CONFIG_DOC) + .define(PULSAR_AUTH_TOKEN_CONFIG, + Type.STRING, + Importance.HIGH, + PULSAR_AUTH_TOKEN_CONFIG_DOC) .define(TOPIC_NAMESPACE_CONFIG, Type.STRING, "public/default", diff --git a/pulsar-io/kafka-connect-adaptor/src/main/java/org/apache/pulsar/io/kafka/connect/PulsarOffsetBackingStore.java b/pulsar-io/kafka-connect-adaptor/src/main/java/org/apache/pulsar/io/kafka/connect/PulsarOffsetBackingStore.java index e616e84f998ac..74b0a6a2eccf8 100644 --- a/pulsar-io/kafka-connect-adaptor/src/main/java/org/apache/pulsar/io/kafka/connect/PulsarOffsetBackingStore.java +++ b/pulsar-io/kafka-connect-adaptor/src/main/java/org/apache/pulsar/io/kafka/connect/PulsarOffsetBackingStore.java @@ -35,6 +35,8 @@ import org.apache.kafka.connect.runtime.WorkerConfig; import org.apache.kafka.connect.storage.OffsetBackingStore; import org.apache.kafka.connect.util.Callback; +import org.apache.pulsar.client.api.AuthenticationFactory; +import org.apache.pulsar.client.api.ClientBuilder; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Producer; @@ -53,6 +55,7 @@ public class PulsarOffsetBackingStore implements OffsetBackingStore { private PulsarClient client; private String serviceUrl; private String topic; + private String token; private Producer producer; private Reader reader; private volatile CompletableFuture outstandingReadToEnd = null; @@ -62,6 +65,7 @@ public void configure(WorkerConfig workerConfig) { this.topic = workerConfig.getString(PulsarKafkaWorkerConfig.OFFSET_STORAGE_TOPIC_CONFIG); checkArgument(!isBlank(topic), "Offset storage topic must be specified"); this.serviceUrl = workerConfig.getString(PulsarKafkaWorkerConfig.PULSAR_SERVICE_URL_CONFIG); + this.token = workerConfig.getString(PulsarKafkaWorkerConfig.PULSAR_AUTH_TOKEN_CONFIG); checkArgument(!isBlank(serviceUrl), "Pulsar service url must be specified at `" + WorkerConfig.BOOTSTRAP_SERVERS_CONFIG + "`"); this.data = new HashMap<>(); @@ -136,9 +140,13 @@ void processMessage(Message message) { @Override public void start() { try { - client = PulsarClient.builder() - .serviceUrl(serviceUrl) - .build(); + ClientBuilder builder = PulsarClient.builder().serviceUrl(serviceUrl); + + if (token != null && token != "") { + builder = builder.authentication(AuthenticationFactory.token(token)); + } + client = builder.build(); + log.info("Successfully created pulsar client to {}", serviceUrl); producer = client.newProducer(Schema.BYTES) .topic(topic)