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