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 54cc1c3  Simplify PushSource (#1898)
54cc1c3 is described below

commit 54cc1c3eb2ed5e69f622e863a7db87a3105a5abb
Author: Sanjeev Kulkarni <[email protected]>
AuthorDate: Sat Jun 2 19:10:14 2018 -0700

    Simplify PushSource (#1898)
    
    ### Motivation
    
    Currently implementing push source requires users implement an additional 
interface called setConsumer and then using that consumer in their loop. This 
pr just exposes a consume method from the pushSource abstract class that 
derived functions can make use of.
---
 .../java/org/apache/pulsar/io/core/PushSource.java     | 18 +++++++-----------
 .../java/org/apache/pulsar/io/kafka/KafkaSource.java   |  9 +--------
 .../org/apache/pulsar/io/rabbitmq/RabbitMQSource.java  | 16 +++++-----------
 .../org/apache/pulsar/io/twitter/TwitterFireHose.java  |  8 +-------
 4 files changed, 14 insertions(+), 37 deletions(-)

diff --git 
a/pulsar-io/core/src/main/java/org/apache/pulsar/io/core/PushSource.java 
b/pulsar-io/core/src/main/java/org/apache/pulsar/io/core/PushSource.java
index 011f8ab..af304b9 100644
--- a/pulsar-io/core/src/main/java/org/apache/pulsar/io/core/PushSource.java
+++ b/pulsar-io/core/src/main/java/org/apache/pulsar/io/core/PushSource.java
@@ -41,16 +41,6 @@ public abstract class PushSource<T> implements Source<T> {
 
     public PushSource() {
         this.queue = new LinkedBlockingQueue<>(this.getQueueLength());
-        this.setConsumer(new Consumer<Record<T>>() {
-            @Override
-            public void accept(Record<T> record) {
-                try {
-                    queue.put(record);
-                } catch (InterruptedException e) {
-                    throw new RuntimeException(e);
-                }
-            }
-        });
     }
 
     @Override
@@ -71,7 +61,13 @@ public abstract class PushSource<T> implements Source<T> {
      * to pass messages whenever there is data to be pushed to Pulsar.
      * @param consumer
      */
-    abstract public void setConsumer(Consumer<Record<T>> consumer);
+    public void consume(Record<T> record) {
+        try {
+            queue.put(record);
+        } catch (InterruptedException e) {
+            throw new RuntimeException(e);
+        }
+    }
 
     /**
      * Get length of the queue that records are push onto
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/KafkaSource.java
index ad967be..9fa43f7 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/KafkaSource.java
@@ -48,8 +48,6 @@ public abstract class KafkaSource<V> extends PushSource<V> {
     private KafkaSourceConfig kafkaSourceConfig;
     Thread runnerThread;
 
-    private java.util.function.Consumer<Record<V>> consumeFunction;
-
     @Override
     public void open(Map<String, Object> config) throws Exception {
         kafkaSourceConfig = KafkaSourceConfig.load(config);
@@ -79,11 +77,6 @@ public abstract class KafkaSource<V> extends PushSource<V> {
     }
 
     @Override
-    public void setConsumer(java.util.function.Consumer<Record<V>> 
consumerFunction) {
-        this.consumeFunction = consumerFunction;
-    }
-
-    @Override
     public void close() throws InterruptedException {
         LOG.info("Stopping kafka source");
         if (runnerThread != null) {
@@ -112,7 +105,7 @@ public abstract class KafkaSource<V> extends PushSource<V> {
                 for (ConsumerRecord<byte[], byte[]> consumerRecord : 
consumerRecords) {
                     LOG.debug("Record received from kafka, key: {}. value: 
{}", consumerRecord.key(), consumerRecord.value());
                     KafkaRecord<V> record = new KafkaRecord<>(consumerRecord, 
extractValue(consumerRecord));
-                    consumeFunction.accept(record);
+                    consume(record);
                     futures[index] = record.getCompletableFuture();
                     index++;
                 }
diff --git 
a/pulsar-io/rabbitmq/src/main/java/org/apache/pulsar/io/rabbitmq/RabbitMQSource.java
 
b/pulsar-io/rabbitmq/src/main/java/org/apache/pulsar/io/rabbitmq/RabbitMQSource.java
index 967c005..e17ab5f 100644
--- 
a/pulsar-io/rabbitmq/src/main/java/org/apache/pulsar/io/rabbitmq/RabbitMQSource.java
+++ 
b/pulsar-io/rabbitmq/src/main/java/org/apache/pulsar/io/rabbitmq/RabbitMQSource.java
@@ -41,7 +41,6 @@ public class RabbitMQSource extends PushSource<byte[]> {
 
     private static Logger logger = 
LoggerFactory.getLogger(RabbitMQSource.class);
 
-    private Consumer<Record<byte[]>> consumer;
     private Connection rabbitMQConnection;
     private Channel rabbitMQChannel;
     private RabbitMQConfig rabbitMQConfig;
@@ -62,33 +61,28 @@ public class RabbitMQSource extends PushSource<byte[]> {
         );
         rabbitMQChannel = rabbitMQConnection.createChannel();
         rabbitMQChannel.queueDeclare(rabbitMQConfig.getQueueName(), false, 
false, false, null);
-        com.rabbitmq.client.Consumer consumer = new 
RabbitMQConsumer(this.consumer, rabbitMQChannel);
+        com.rabbitmq.client.Consumer consumer = new RabbitMQConsumer(this, 
rabbitMQChannel);
         rabbitMQChannel.basicConsume(rabbitMQConfig.getQueueName(), consumer);
         logger.info("A consumer for queue {} has been successfully started.", 
rabbitMQConfig.getQueueName());
     }
 
     @Override
-    public void setConsumer(Consumer<Record<byte[]>> consumer) {
-        this.consumer = consumer;
-    }
-
-    @Override
     public void close() throws Exception {
         rabbitMQChannel.close();
         rabbitMQConnection.close();
     }
 
     private class RabbitMQConsumer extends DefaultConsumer {
-        private Consumer<Record<byte[]>> consumeFunction;
+        private RabbitMQSource source;
 
-        public RabbitMQConsumer(Consumer<Record<byte[]>> consumeFunction, 
Channel channel) {
+        public RabbitMQConsumer(RabbitMQSource source, Channel channel) {
             super(channel);
-            this.consumeFunction = consumeFunction;
+            this.source = source;
         }
 
         @Override
         public void handleDelivery(String consumerTag, Envelope envelope, 
AMQP.BasicProperties properties, byte[] body) throws IOException {
-            consumeFunction.accept(new RabbitMQRecord(body));
+            source.consume(new RabbitMQRecord(body));
         }
     }
 
diff --git 
a/pulsar-io/twitter/src/main/java/org/apache/pulsar/io/twitter/TwitterFireHose.java
 
b/pulsar-io/twitter/src/main/java/org/apache/pulsar/io/twitter/TwitterFireHose.java
index 8331e7d..3ecfb09 100644
--- 
a/pulsar-io/twitter/src/main/java/org/apache/pulsar/io/twitter/TwitterFireHose.java
+++ 
b/pulsar-io/twitter/src/main/java/org/apache/pulsar/io/twitter/TwitterFireHose.java
@@ -50,7 +50,6 @@ public class TwitterFireHose extends PushSource<String> {
 
     // ----- Runtime fields
     private Object waitObject;
-    private Consumer<Record<String>> consumeFunction;
 
     @Override
     public void open(Map<String, Object> config) throws IOException {
@@ -66,11 +65,6 @@ public class TwitterFireHose extends PushSource<String> {
     }
 
     @Override
-    public void setConsumer(Consumer<Record<String>> consumeFunction) {
-        this.consumeFunction = consumeFunction;
-    }
-
-    @Override
     public void close() throws Exception {
         stopThread();
     }
@@ -125,7 +119,7 @@ public class TwitterFireHose extends PushSource<String> {
                             // We don't really care if the record succeeds or 
not.
                             // However might be in the future to count failures
                             // TODO:- Figure out the metrics story for 
connectors
-                            consumeFunction.accept(new TwitterRecord(line));
+                            consume(new TwitterRecord(line));
                         } catch (Exception e) {
                             LOG.error("Exception thrown");
                         }

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

Reply via email to