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