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].