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;
}
}