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

jerrypeng 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 577e579  When fetching from kafka, specify to get the earliest of 
messages (#1896)
577e579 is described below

commit 577e5793abbb8d2bd1ce0295cc88d84fc1acf3db
Author: Sanjeev Kulkarni <[email protected]>
AuthorDate: Sat Jun 2 12:34:38 2018 -0700

    When fetching from kafka, specify to get the earliest of messages (#1896)
---
 .../kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaSource.java      | 1 +
 1 file changed, 1 insertion(+)

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 fa73418..ad967be 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
@@ -69,6 +69,7 @@ public abstract class KafkaSource<V> extends PushSource<V> {
         props.put(ConsumerConfig.FETCH_MIN_BYTES_CONFIG, 
kafkaSourceConfig.getFetchMinBytes().toString());
         props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 
kafkaSourceConfig.getAutoCommitIntervalMs().toString());
         props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 
kafkaSourceConfig.getSessionTimeoutMs().toString());
+        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
 
         props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, 
kafkaSourceConfig.getKeyDeserializationClass());
         props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, 
kafkaSourceConfig.getValueDeserializationClass());

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

Reply via email to