sijie closed pull request #1898: Simplify PushSource
URL: https://github.com/apache/incubator-pulsar/pull/1898
This is a PR merged from a forked repository.
As GitHub hides the original diff on merge, it is displayed below for
the sake of provenance:
As this is a foreign pull request (from a fork), the diff is supplied
below (as it won't show otherwise due to GitHub magic):
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 011f8ab386..af304b937e 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 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 void accept(Record<T> record) {
* 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 ad967be2cb..9fa43f7ca3 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 @@
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);
@@ -78,11 +76,6 @@ public void open(Map<String, Object> config) throws
Exception {
}
- @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");
@@ -112,7 +105,7 @@ public void start() {
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 967c00530e..e17ab5f8d7 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 @@
private static Logger logger =
LoggerFactory.getLogger(RabbitMQSource.class);
- private Consumer<Record<byte[]>> consumer;
private Connection rabbitMQConnection;
private Channel rabbitMQChannel;
private RabbitMQConfig rabbitMQConfig;
@@ -62,16 +61,11 @@ public void open(Map<String, Object> config) throws
Exception {
);
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();
@@ -79,16 +73,16 @@ public void close() throws Exception {
}
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 8331e7d3d8..3ecfb090bf 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 @@
// ----- Runtime fields
private Object waitObject;
- private Consumer<Record<String>> consumeFunction;
@Override
public void open(Map<String, Object> config) throws IOException {
@@ -65,11 +64,6 @@ public void open(Map<String, Object> config) throws
IOException {
startThread(hoseConfig);
}
- @Override
- public void setConsumer(Consumer<Record<String>> consumeFunction) {
- this.consumeFunction = consumeFunction;
- }
-
@Override
public void close() throws Exception {
stopThread();
@@ -125,7 +119,7 @@ public boolean process() throws IOException,
InterruptedException {
// 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");
}
----------------------------------------------------------------
This is an automated message from the Apache Git Service.
To respond to the message, please log on GitHub and use the
URL above to go to the specific comment.
For queries about this service, please contact Infrastructure at:
[email protected]
With regards,
Apache Git Services