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

sijie pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-pulsar.git


The following commit(s) were added to refs/heads/master by this push:
     new 1950538  Added a bunch of concrete sources/sinks so that they are 
usable without having to write code (#1934)
1950538 is described below

commit 19505381e4cb0ddbaaab6d2c76de2b8fd8a0f10d
Author: Sanjeev Kulkarni <[email protected]>
AuthorDate: Fri Jun 8 09:14:21 2018 -0700

    Added a bunch of concrete sources/sinks so that they are usable without 
having to write code (#1934)
    
    Currently Kafka/Cassandra/Aerospike connectors are all abstract which means 
that to use them, one has to actually write a concrete class. This pr adds a 
bunch of concrete classes so that these classes are usable right away without 
having to write any code
---
 ...rospikeSink.java => AerospikeAbstractSink.java} |  4 +--
 .../pulsar/io/aerospike/AerospikeStringSink.java   | 33 ++++++++++++++++++++++
 ...ssandraSink.java => CassandraAbstractSink.java} |  4 +--
 .../pulsar/io/cassandra/CassandraStringSink.java   | 33 ++++++++++++++++++++++
 .../{KafkaSink.java => KafkaAbstractSink.java}     |  4 +--
 .../{KafkaSource.java => KafkaAbstractSource.java} |  4 +--
 .../apache/pulsar/io/kafka/KafkaStringSink.java    | 33 ++++++++++++++++++++++
 .../apache/pulsar/io/kafka/KafkaStringSource.java  | 33 ++++++++++++++++++++++
 8 files changed, 140 insertions(+), 8 deletions(-)

diff --git 
a/pulsar-io/aerospike/src/main/java/org/apache/pulsar/io/aerospike/AerospikeSink.java
 
b/pulsar-io/aerospike/src/main/java/org/apache/pulsar/io/aerospike/AerospikeAbstractSink.java
similarity index 98%
rename from 
pulsar-io/aerospike/src/main/java/org/apache/pulsar/io/aerospike/AerospikeSink.java
rename to 
pulsar-io/aerospike/src/main/java/org/apache/pulsar/io/aerospike/AerospikeAbstractSink.java
index 18e8519..6325eec 100644
--- 
a/pulsar-io/aerospike/src/main/java/org/apache/pulsar/io/aerospike/AerospikeSink.java
+++ 
b/pulsar-io/aerospike/src/main/java/org/apache/pulsar/io/aerospike/AerospikeAbstractSink.java
@@ -45,9 +45,9 @@ import java.util.concurrent.LinkedBlockingDeque;
  * A Simple abstract class for Aerospike sink
  * Users need to implement extractKeyValue function to use this sink
  */
-public abstract class AerospikeSink<K, V> extends SimpleSink<byte[]> {
+public abstract class AerospikeAbstractSink<K, V> extends SimpleSink<byte[]> {
 
-    private static final Logger LOG = 
LoggerFactory.getLogger(AerospikeSink.class);
+    private static final Logger LOG = 
LoggerFactory.getLogger(AerospikeAbstractSink.class);
 
     // ----- Runtime fields
     private AerospikeSinkConfig aerospikeSinkConfig;
diff --git 
a/pulsar-io/aerospike/src/main/java/org/apache/pulsar/io/aerospike/AerospikeStringSink.java
 
b/pulsar-io/aerospike/src/main/java/org/apache/pulsar/io/aerospike/AerospikeStringSink.java
new file mode 100644
index 0000000..affaa84
--- /dev/null
+++ 
b/pulsar-io/aerospike/src/main/java/org/apache/pulsar/io/aerospike/AerospikeStringSink.java
@@ -0,0 +1,33 @@
+/**
+ * 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.pulsar.io.aerospike;
+
+import org.apache.pulsar.common.util.KeyValue;
+
+/**
+ * Aerospike sink that treats incoming messages on the input topic as Strings
+ * and write identical key/value pairs.
+ */
+public class AerospikeStringSink extends AerospikeAbstractSink<String, String> 
{
+    @Override
+    public KeyValue<String, String> extractKeyValue(byte[] message) {
+        return new KeyValue<>(new String(message), new String(message));
+    }
+}
\ No newline at end of file
diff --git 
a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java
 
b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraAbstractSink.java
similarity index 97%
rename from 
pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java
rename to 
pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraAbstractSink.java
index 710feb7..c177f9f 100644
--- 
a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java
+++ 
b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraAbstractSink.java
@@ -39,9 +39,9 @@ import java.util.concurrent.CompletableFuture;
  * A Simple abstract class for Cassandra sink
  * Users need to implement extractKeyValue function to use this sink
  */
-public abstract class CassandraSink<K, V> extends SimpleSink<byte[]> {
+public abstract class CassandraAbstractSink<K, V> extends SimpleSink<byte[]> {
 
-    private static final Logger LOG = 
LoggerFactory.getLogger(CassandraSink.class);
+    private static final Logger LOG = 
LoggerFactory.getLogger(CassandraAbstractSink.class);
 
     // ----- Runtime fields
     private Cluster cluster;
diff --git 
a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java
 
b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java
new file mode 100644
index 0000000..c4026d2
--- /dev/null
+++ 
b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java
@@ -0,0 +1,33 @@
+/**
+ * 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.pulsar.io.cassandra;
+
+import org.apache.pulsar.common.util.KeyValue;
+
+/**
+ * Cassandra sink that treats incoming messages on the input topic as Strings
+ * and write identical key/value pairs.
+ */
+public class CassandraStringSink extends CassandraAbstractSink<String, String> 
{
+    @Override
+    public KeyValue<String, String> extractKeyValue(byte[] message) {
+        return new KeyValue<>(new String(message), new String(message));
+    }
+}
\ No newline at end of file
diff --git 
a/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaSink.java 
b/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaAbstractSink.java
similarity index 97%
rename from 
pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaSink.java
rename to 
pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaAbstractSink.java
index 5081d54..ae81693 100644
--- a/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaSink.java
+++ 
b/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaAbstractSink.java
@@ -39,9 +39,9 @@ import java.util.concurrent.Future;
  * A Simple abstract class for Kafka sink
  * Users need to implement extractKeyValue function to use this sink
  */
-public abstract class KafkaSink<K, V> extends SimpleSink<byte[]> {
+public abstract class KafkaAbstractSink<K, V> extends SimpleSink<byte[]> {
 
-    private static final Logger LOG = LoggerFactory.getLogger(KafkaSink.class);
+    private static final Logger LOG = 
LoggerFactory.getLogger(KafkaAbstractSink.class);
 
     private Producer<K, V> producer;
     private Properties props = new Properties();
diff --git 
a/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaSource.java 
b/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaAbstractSource.java
similarity index 98%
rename from 
pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaSource.java
rename to 
pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaAbstractSource.java
index 9fa43f7..05d90b3 100644
--- a/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaSource.java
+++ 
b/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaAbstractSource.java
@@ -39,9 +39,9 @@ import java.util.concurrent.ExecutionException;
 /**
  * Simple Kafka Source to transfer messages from a Kafka topic
  */
-public abstract class KafkaSource<V> extends PushSource<V> {
+public abstract class KafkaAbstractSource<V> extends PushSource<V> {
 
-    private static final Logger LOG = 
LoggerFactory.getLogger(KafkaSource.class);
+    private static final Logger LOG = 
LoggerFactory.getLogger(KafkaAbstractSource.class);
 
     private Consumer<byte[], byte[]> consumer;
     private Properties props;
diff --git 
a/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaStringSink.java 
b/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaStringSink.java
new file mode 100644
index 0000000..6e39817
--- /dev/null
+++ 
b/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaStringSink.java
@@ -0,0 +1,33 @@
+/**
+ * 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.pulsar.io.kafka;
+
+import org.apache.pulsar.common.util.KeyValue;
+
+/**
+ * Kafka sink that treats incoming messages on the input topic as Strings
+ * and write identical key/value pairs.
+ */
+public class KafkaStringSink extends KafkaAbstractSink<String, String> {
+    @Override
+    public KeyValue<String, String> extractKeyValue(byte[] message) {
+        return new KeyValue<>(new String(message), new String(message));
+    }
+}
\ No newline at end of file
diff --git 
a/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaStringSource.java
 
b/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaStringSource.java
new file mode 100644
index 0000000..e31b75f
--- /dev/null
+++ 
b/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaStringSource.java
@@ -0,0 +1,33 @@
+/**
+ * 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.pulsar.io.kafka;
+
+import org.apache.kafka.clients.consumer.*;
+
+/**
+ * Simple Kafka Source that just transfers the value part of the kafka records
+ * as Strings
+ */
+public class KafkaStringSource extends KafkaAbstractSource<String> {
+    @Override
+    public String extractValue(ConsumerRecord<byte[], byte[]> record) {
+        return new String(record.value());
+    }
+}

-- 
To stop receiving notification emails like this one, please contact
[email protected].

Reply via email to