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

riemer pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git


The following commit(s) were added to refs/heads/dev by this push:
     new 58d219491 fix: kafka consumer data loss promble (#1629)
58d219491 is described below

commit 58d2194912d1d7e86faf63a85c6c3a42dc50e740
Author: luoluoyuyu <[email protected]>
AuthorDate: Tue Jun 27 02:52:49 2023 +0800

    fix: kafka consumer data loss promble (#1629)
    
    * add sample configuration of pulsar subscription-name
    
    * fix: Kafka consumer data loss promble
    
    * add sample configuration of pulsar subscription-name
    
    * kafka adapter adds AUTO_OFFSET_RESET configuration option
    
    * kafka adapter adds AUTO_OFFSET_RESET configuration option
    
    * fix:add default configuration items
---
 .../iiot/protocol/stream/KafkaProtocol.java        | 21 +++++++++--
 .../strings.en                                     | 12 ++++++
 .../pe/shared/config/kafka/KafkaConfig.java        | 15 +++++++-
 .../pe/shared/config/kafka/KafkaConnectUtils.java  | 32 +++++++++++++++-
 .../pe/shared/config/kafka/kafka/KafkaConfig.java  | 13 ++++++-
 .../config/kafka/kafka/KafkaConnectUtils.java      | 35 +++++++++++++++++-
 .../kafka/config/AutoOffsetResetConfig.java        | 43 ++++++++++++++++++++++
 .../kafka/config/ConsumerConfigFactory.java        |  3 +-
 8 files changed, 164 insertions(+), 10 deletions(-)

diff --git 
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/KafkaProtocol.java
 
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/KafkaProtocol.java
index 9f5ce94f3..e65806cb1 100644
--- 
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/KafkaProtocol.java
+++ 
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/KafkaProtocol.java
@@ -33,6 +33,7 @@ import 
org.apache.streampipes.extensions.api.extractor.IStaticPropertyExtractor;
 import org.apache.streampipes.extensions.api.runtime.SupportsRuntimeConfig;
 import 
org.apache.streampipes.extensions.management.connect.adapter.parser.Parsers;
 import org.apache.streampipes.messaging.kafka.SpKafkaConsumer;
+import org.apache.streampipes.messaging.kafka.config.KafkaConfigAppender;
 import org.apache.streampipes.model.AdapterType;
 import org.apache.streampipes.model.connect.guess.GuessSchema;
 import org.apache.streampipes.model.grounding.KafkaTransportProtocol;
@@ -40,10 +41,10 @@ import 
org.apache.streampipes.model.grounding.SimpleTopicDefinition;
 import org.apache.streampipes.model.staticproperty.Option;
 import 
org.apache.streampipes.model.staticproperty.RuntimeResolvableOneOfStaticProperty;
 import org.apache.streampipes.model.staticproperty.StaticProperty;
+import org.apache.streampipes.model.staticproperty.StaticPropertyAlternative;
 import org.apache.streampipes.pe.shared.config.kafka.kafka.KafkaConfig;
 import org.apache.streampipes.pe.shared.config.kafka.kafka.KafkaConnectUtils;
 import org.apache.streampipes.sdk.builder.adapter.AdapterConfigurationBuilder;
-import org.apache.streampipes.sdk.helpers.Labels;
 import org.apache.streampipes.sdk.helpers.Locales;
 import org.apache.streampipes.sdk.utils.Assets;
 
@@ -89,6 +90,7 @@ public class KafkaProtocol implements StreamPipesAdapter, 
SupportsRuntimeConfig
     final Properties props = new Properties();
 
     kafkaConfig.getSecurityConfig().appendConfig(props);
+    kafkaConfig.getAutoOffsetResetConfig().appendConfig(props);
 
     props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,
         kafkaConfig.getKafkaHost() + ":" + kafkaConfig.getKafkaPort());
@@ -137,6 +139,10 @@ public class KafkaProtocol implements StreamPipesAdapter, 
SupportsRuntimeConfig
 
   @Override
   public IAdapterConfiguration declareConfig() {
+
+    StaticPropertyAlternative latestAlternative = 
KafkaConnectUtils.getAlternativesLatest();
+    latestAlternative.setSelected(true);
+
     return AdapterConfigurationBuilder
         .create(ID, KafkaProtocol::new)
         .withSupportedParsers(Parsers.defaultParsers())
@@ -144,7 +150,7 @@ public class KafkaProtocol implements StreamPipesAdapter, 
SupportsRuntimeConfig
         .withLocales(Locales.EN)
         .withCategory(AdapterType.Generic, AdapterType.Manufacturing)
 
-        .requiredAlternatives(Labels.withId(KafkaConnectUtils.ACCESS_MODE),
+        .requiredAlternatives(KafkaConnectUtils.getAccessModeLabel(),
             KafkaConnectUtils.getAlternativeUnauthenticatedPlain(),
             KafkaConnectUtils.getAlternativeUnauthenticatedSSL(),
             KafkaConnectUtils.getAlternativesSaslPlain(),
@@ -158,6 +164,10 @@ public class KafkaProtocol implements StreamPipesAdapter, 
SupportsRuntimeConfig
         
.requiredSingleValueSelectionFromContainer(KafkaConnectUtils.getTopicLabel(), 
Arrays.asList(
             KafkaConnectUtils.HOST_KEY,
             KafkaConnectUtils.PORT_KEY))
+        
.requiredAlternatives(KafkaConnectUtils.getAutoOffsetResetConfigLabel(),
+                KafkaConnectUtils.getAlternativesEarliest(),
+                latestAlternative,
+                KafkaConnectUtils.getAlternativesNone())
         .buildConfiguration();
   }
 
@@ -171,12 +181,17 @@ public class KafkaProtocol implements StreamPipesAdapter, 
SupportsRuntimeConfig
     protocol.setBrokerHostname(config.getKafkaHost());
     protocol.setTopicDefinition(new SimpleTopicDefinition(config.getTopic()));
 
+    List<KafkaConfigAppender> kafkaConfigAppenderList = new ArrayList<>(2);
+    kafkaConfigAppenderList.add(this.config.getSecurityConfig());
+    kafkaConfigAppenderList.add(this.config.getAutoOffsetResetConfig());
+
     this.kafkaConsumer = new SpKafkaConsumer(protocol,
         config.getTopic(),
         new BrokerEventProcessor(extractor.selectedParser(), (event) -> {
           collector.collect(event);
         }),
-        Collections.singletonList(this.config.getSecurityConfig()));
+        kafkaConfigAppenderList
+        );
 
     thread = new Thread(this.kafkaConsumer);
     thread.start();
diff --git 
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/resources/org.apache.streampipes.connect.iiot.protocol.stream.kafka/strings.en
 
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/resources/org.apache.streampipes.connect.iiot.protocol.stream.kafka/strings.en
index 15d3dce84..fd8ef54bc 100644
--- 
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/resources/org.apache.streampipes.connect.iiot.protocol.stream.kafka/strings.en
+++ 
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/resources/org.apache.streampipes.connect.iiot.protocol.stream.kafka/strings.en
@@ -59,3 +59,15 @@ value-deserialization.description=
 
 hide-internal-topics.title=Hide internal topics
 hide-internal-topics.description=Do not show topics that are only used 
internally
+
+auto.offset.reset.title=Auto Offset Reset
+auto.offset.reset.description=Configure the starting offset for consumer 
consumption when there is no initial offset in Kafka or when the current offset 
does not exist
+
+earliest.title=Earliest
+earliest.description=Offsets are initialized to the earliest
+
+latest.title=Latest
+latest.description=Offsets are initialized to the Latest
+
+none.title=None
+none.description=Consumer throws exceptions
\ No newline at end of file
diff --git 
a/streampipes-extensions/streampipes-pipeline-elements-shared/src/main/java/org/apache/streampipes/pe/shared/config/kafka/KafkaConfig.java
 
b/streampipes-extensions/streampipes-pipeline-elements-shared/src/main/java/org/apache/streampipes/pe/shared/config/kafka/KafkaConfig.java
index 16eb4016d..7e1cfecaf 100644
--- 
a/streampipes-extensions/streampipes-pipeline-elements-shared/src/main/java/org/apache/streampipes/pe/shared/config/kafka/KafkaConfig.java
+++ 
b/streampipes-extensions/streampipes-pipeline-elements-shared/src/main/java/org/apache/streampipes/pe/shared/config/kafka/KafkaConfig.java
@@ -18,6 +18,7 @@
 
 package org.apache.streampipes.pe.shared.config.kafka;
 
+import org.apache.streampipes.messaging.kafka.config.AutoOffsetResetConfig;
 import org.apache.streampipes.messaging.kafka.security.KafkaSecurityConfig;
 
 public class KafkaConfig {
@@ -27,16 +28,18 @@ public class KafkaConfig {
   private String topic;
 
   KafkaSecurityConfig securityConfig;
-
+  AutoOffsetResetConfig autoOffsetResetConfig;
 
   public KafkaConfig(String kafkaHost,
                      Integer kafkaPort,
                      String topic,
-                     KafkaSecurityConfig securityConfig) {
+                     KafkaSecurityConfig securityConfig,
+                     AutoOffsetResetConfig autoOffsetResetConfig) {
     this.kafkaHost = kafkaHost;
     this.kafkaPort = kafkaPort;
     this.topic = topic;
     this.securityConfig = securityConfig;
+    this.autoOffsetResetConfig = autoOffsetResetConfig;
   }
 
   public String getKafkaHost() {
@@ -71,4 +74,12 @@ public class KafkaConfig {
     this.securityConfig = securityConfig;
   }
 
+  public void setAutoOffsetResetConfig(AutoOffsetResetConfig 
autoOffsetResetConfig) {
+    this.autoOffsetResetConfig = autoOffsetResetConfig;
+  }
+
+  public AutoOffsetResetConfig getAutoOffsetResetConfig() {
+    return autoOffsetResetConfig;
+  }
+
 }
diff --git 
a/streampipes-extensions/streampipes-pipeline-elements-shared/src/main/java/org/apache/streampipes/pe/shared/config/kafka/KafkaConnectUtils.java
 
b/streampipes-extensions/streampipes-pipeline-elements-shared/src/main/java/org/apache/streampipes/pe/shared/config/kafka/KafkaConnectUtils.java
index 124c1c6e8..7cc3f2430 100644
--- 
a/streampipes-extensions/streampipes-pipeline-elements-shared/src/main/java/org/apache/streampipes/pe/shared/config/kafka/KafkaConnectUtils.java
+++ 
b/streampipes-extensions/streampipes-pipeline-elements-shared/src/main/java/org/apache/streampipes/pe/shared/config/kafka/KafkaConnectUtils.java
@@ -18,6 +18,9 @@
 
 package org.apache.streampipes.pe.shared.config.kafka;
 
+
+
+import org.apache.streampipes.messaging.kafka.config.AutoOffsetResetConfig;
 import org.apache.streampipes.messaging.kafka.security.KafkaSecurityConfig;
 import 
org.apache.streampipes.messaging.kafka.security.KafkaSecuritySaslPlainConfig;
 import 
org.apache.streampipes.messaging.kafka.security.KafkaSecuritySaslSSLConfig;
@@ -30,6 +33,8 @@ import org.apache.streampipes.sdk.helpers.Alternatives;
 import org.apache.streampipes.sdk.helpers.Label;
 import org.apache.streampipes.sdk.helpers.Labels;
 
+import org.apache.kafka.clients.consumer.ConsumerConfig;
+
 public class KafkaConnectUtils {
 
   public static final String TOPIC_KEY = "topic";
@@ -49,6 +54,11 @@ public class KafkaConnectUtils {
 
   private static final String HIDE_INTERNAL_TOPICS = "hide-internal-topics";
 
+  public static final String AUTO_OFFSET_RESET_CONFIG = 
ConsumerConfig.AUTO_OFFSET_RESET_CONFIG;
+  public static final String EARLIEST = "earliest";
+  public static final String LATEST = "latest";
+  public static final String NONE = "none";
+
   public static Label getTopicLabel() {
     return Labels.withId(TOPIC_KEY);
   }
@@ -73,6 +83,10 @@ public class KafkaConnectUtils {
     return Labels.withId(ACCESS_MODE);
   }
 
+  public static Label getAutoOffsetResetConfigLabel() {
+    return Labels.withId(AUTO_OFFSET_RESET_CONFIG);
+  }
+
   public static KafkaConfig getConfig(StaticPropertyExtractor extractor, 
boolean containsTopic) {
     String brokerUrl = extractor.singleValueParameter(HOST_KEY, String.class);
     String topic = "";
@@ -102,7 +116,10 @@ public class KafkaConnectUtils {
           new KafkaSecurityUnauthenticatedPlainConfig();
     }
 
-    return new KafkaConfig(brokerUrl, port, topic, securityConfig);
+    String auto = 
extractor.selectedAlternativeInternalId(AUTO_OFFSET_RESET_CONFIG);
+    AutoOffsetResetConfig autoOffsetResetConfig = new 
AutoOffsetResetConfig(auto);
+
+    return new KafkaConfig(brokerUrl, port, topic, securityConfig, 
autoOffsetResetConfig);
   }
 
   private static boolean isUseSSL(String authentication) {
@@ -136,4 +153,17 @@ public class KafkaConnectUtils {
             
StaticProperties.stringFreeTextProperty(Labels.withId(KafkaConnectUtils.USERNAME_KEY)),
             
StaticProperties.secretValue(Labels.withId(KafkaConnectUtils.PASSWORD_KEY))));
   }
+
+
+  public static StaticPropertyAlternative getAlternativesLatest() {
+    return Alternatives.from(Labels.withId(LATEST));
+  }
+
+  public static StaticPropertyAlternative getAlternativesEarliest() {
+    return Alternatives.from(Labels.withId(EARLIEST));
+  }
+
+  public static StaticPropertyAlternative getAlternativesNone() {
+    return Alternatives.from(Labels.withId(NONE));
+  }
 }
diff --git 
a/streampipes-extensions/streampipes-pipeline-elements-shared/src/main/java/org/apache/streampipes/pe/shared/config/kafka/kafka/KafkaConfig.java
 
b/streampipes-extensions/streampipes-pipeline-elements-shared/src/main/java/org/apache/streampipes/pe/shared/config/kafka/kafka/KafkaConfig.java
index 1a66e0fa1..61c83c2a1 100644
--- 
a/streampipes-extensions/streampipes-pipeline-elements-shared/src/main/java/org/apache/streampipes/pe/shared/config/kafka/kafka/KafkaConfig.java
+++ 
b/streampipes-extensions/streampipes-pipeline-elements-shared/src/main/java/org/apache/streampipes/pe/shared/config/kafka/kafka/KafkaConfig.java
@@ -18,6 +18,7 @@
 
 package org.apache.streampipes.pe.shared.config.kafka.kafka;
 
+import org.apache.streampipes.messaging.kafka.config.AutoOffsetResetConfig;
 import org.apache.streampipes.messaging.kafka.security.KafkaSecurityConfig;
 
 public class KafkaConfig {
@@ -27,15 +28,18 @@ public class KafkaConfig {
   private String topic;
 
   KafkaSecurityConfig securityConfig;
+  AutoOffsetResetConfig autoOffsetResetConfig;
 
   public KafkaConfig(String kafkaHost,
                      Integer kafkaPort,
                      String topic,
-                     KafkaSecurityConfig securityConfig) {
+                     KafkaSecurityConfig securityConfig,
+                     AutoOffsetResetConfig autoOffsetResetConfig) {
     this.kafkaHost = kafkaHost;
     this.kafkaPort = kafkaPort;
     this.topic = topic;
     this.securityConfig = securityConfig;
+    this.autoOffsetResetConfig = autoOffsetResetConfig;
   }
 
   public String getKafkaHost() {
@@ -70,4 +74,11 @@ public class KafkaConfig {
     this.securityConfig = securityConfig;
   }
 
+  public AutoOffsetResetConfig getAutoOffsetResetConfig() {
+    return autoOffsetResetConfig;
+  }
+
+  public void setAutoOffsetResetConfig(AutoOffsetResetConfig 
autoOffsetResetConfig) {
+    this.autoOffsetResetConfig = autoOffsetResetConfig;
+  }
 }
diff --git 
a/streampipes-extensions/streampipes-pipeline-elements-shared/src/main/java/org/apache/streampipes/pe/shared/config/kafka/kafka/KafkaConnectUtils.java
 
b/streampipes-extensions/streampipes-pipeline-elements-shared/src/main/java/org/apache/streampipes/pe/shared/config/kafka/kafka/KafkaConnectUtils.java
index cc9bfbff1..6ad3f61db 100644
--- 
a/streampipes-extensions/streampipes-pipeline-elements-shared/src/main/java/org/apache/streampipes/pe/shared/config/kafka/kafka/KafkaConnectUtils.java
+++ 
b/streampipes-extensions/streampipes-pipeline-elements-shared/src/main/java/org/apache/streampipes/pe/shared/config/kafka/kafka/KafkaConnectUtils.java
@@ -19,6 +19,7 @@
 package org.apache.streampipes.pe.shared.config.kafka.kafka;
 
 import 
org.apache.streampipes.extensions.api.extractor.IStaticPropertyExtractor;
+import org.apache.streampipes.messaging.kafka.config.AutoOffsetResetConfig;
 import org.apache.streampipes.messaging.kafka.security.KafkaSecurityConfig;
 import 
org.apache.streampipes.messaging.kafka.security.KafkaSecuritySaslPlainConfig;
 import 
org.apache.streampipes.messaging.kafka.security.KafkaSecuritySaslSSLConfig;
@@ -30,6 +31,8 @@ import org.apache.streampipes.sdk.helpers.Alternatives;
 import org.apache.streampipes.sdk.helpers.Label;
 import org.apache.streampipes.sdk.helpers.Labels;
 
+import org.apache.kafka.clients.consumer.ConsumerConfig;
+
 public class KafkaConnectUtils {
 
   public static final String TOPIC_KEY = "topic";
@@ -55,6 +58,11 @@ public class KafkaConnectUtils {
 
   private static final String HIDE_INTERNAL_TOPICS = "hide-internal-topics";
 
+  public static final String AUTO_OFFSET_RESET_CONFIG = 
ConsumerConfig.AUTO_OFFSET_RESET_CONFIG;
+  public static final String EARLIEST = "earliest";
+  public static final String LATEST = "latest";
+  public static final String NONE = "none";
+
   public static Label getTopicLabel() {
     return Labels.withId(TOPIC_KEY);
   }
@@ -79,6 +87,11 @@ public class KafkaConnectUtils {
     return Labels.withId(ACCESS_MODE);
   }
 
+  public static Label getAutoOffsetResetConfigLabel() {
+    return Labels.withId(AUTO_OFFSET_RESET_CONFIG);
+  }
+
+
   public static KafkaConfig getConfig(IStaticPropertyExtractor extractor, 
boolean containsTopic) {
     String brokerUrl = extractor.singleValueParameter(HOST_KEY, String.class);
     String topic = "";
@@ -93,7 +106,7 @@ public class KafkaConnectUtils {
 
     KafkaSecurityConfig securityConfig;
 
-    //KafkaSerializerConfig serializerConfig = new 
KafkaSerializerByteArrayConfig();
+    //KafkaSerializerConfig serializerConfig = new 
KafkaSerializerByteArrayConfig()
 
     // check if a user for the authentication is defined
     if (authentication.equals(KafkaConnectUtils.SASL_SSL) || 
authentication.equals(KafkaConnectUtils.SASL_PLAIN)) {
@@ -110,7 +123,12 @@ public class KafkaConnectUtils {
           new KafkaSecurityUnauthenticatedPlainConfig();
     }
 
-    return new KafkaConfig(brokerUrl, port, topic, securityConfig);
+
+
+    String auto = 
extractor.selectedAlternativeInternalId(AUTO_OFFSET_RESET_CONFIG);
+    AutoOffsetResetConfig autoOffsetResetConfig = new 
AutoOffsetResetConfig(auto);
+
+    return new KafkaConfig(brokerUrl, port, topic, securityConfig, 
autoOffsetResetConfig);
   }
 
   private static boolean isUseSSL(String authentication) {
@@ -144,4 +162,17 @@ public class KafkaConnectUtils {
             
StaticProperties.stringFreeTextProperty(Labels.withId(KafkaConnectUtils.USERNAME_KEY)),
             
StaticProperties.secretValue(Labels.withId(KafkaConnectUtils.PASSWORD_KEY))));
   }
+
+
+  public static StaticPropertyAlternative getAlternativesLatest() {
+    return Alternatives.from(Labels.withId(LATEST));
+  }
+
+  public static StaticPropertyAlternative getAlternativesEarliest() {
+    return Alternatives.from(Labels.withId(EARLIEST));
+  }
+
+  public static StaticPropertyAlternative getAlternativesNone() {
+    return Alternatives.from(Labels.withId(NONE));
+  }
 }
diff --git 
a/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/config/AutoOffsetResetConfig.java
 
b/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/config/AutoOffsetResetConfig.java
new file mode 100644
index 000000000..25f8ae29f
--- /dev/null
+++ 
b/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/config/AutoOffsetResetConfig.java
@@ -0,0 +1,43 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.streampipes.messaging.kafka.config;
+
+import org.apache.kafka.clients.consumer.ConsumerConfig;
+
+import java.util.Properties;
+
+public class AutoOffsetResetConfig implements KafkaConfigAppender {
+
+  private final String autoOffsetResetConfig;
+
+  @Override
+  public void appendConfig(Properties props) {
+    props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, autoOffsetResetConfig);
+  }
+
+  public AutoOffsetResetConfig(String autoOffsetResetConfig) {
+    this.autoOffsetResetConfig = autoOffsetResetConfig;
+  }
+
+
+  public String getAutoOffsetResetConfig() {
+    return autoOffsetResetConfig;
+  }
+
+}
diff --git 
a/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/config/ConsumerConfigFactory.java
 
b/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/config/ConsumerConfigFactory.java
index 461d96e04..4019cd7ba 100644
--- 
a/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/config/ConsumerConfigFactory.java
+++ 
b/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/config/ConsumerConfigFactory.java
@@ -33,6 +33,7 @@ public class ConsumerConfigFactory extends 
AbstractConfigFactory {
   private static final Integer FETCH_MAX_BYTES_CONFIG_DEFAULT = 52428800;
   private static final String KEY_DESERIALIZER_CLASS_CONFIG_DEFAULT = 
ByteArrayDeserializer.class.getName();
   private static final String VALUE_DESERIALIZER_CLASS_CONFIG_DEFAULT = 
ByteArrayDeserializer.class.getName();
+  private static final String AUTO_OFFSET_RESET_CONFIG_DEFAULT = "earliest";
 
   public ConsumerConfigFactory(KafkaTransportProtocol protocol) {
     super(protocol);
@@ -55,7 +56,7 @@ public class ConsumerConfigFactory extends 
AbstractConfigFactory {
 
 
     props.put(ConsumerConfig.CLIENT_ID_CONFIG, UUID.randomUUID().toString());
-
+    props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, 
AUTO_OFFSET_RESET_CONFIG_DEFAULT);
     return props;
   }
 }

Reply via email to