This is an automated email from the ASF dual-hosted git repository.
markli pushed a commit to branch kafka-parse-error
in repository https://gitbox.apache.org/repos/asf/incubator-eventmesh.git
The following commit(s) were added to refs/heads/kafka-parse-error by this push:
new 50522e67c Updates on Key and Value de/serializer
50522e67c is described below
commit 50522e67c8ee930ba61cd5d9b10641ca12b5e18a
Author: Mark Li <[email protected]>
AuthorDate: Mon Feb 27 19:21:12 2023 -0800
Updates on Key and Value de/serializer
---
.../eventmesh/connector/kafka/consumer/ConsumerImpl.java | 10 +++++-----
.../eventmesh/connector/kafka/producer/ProducerImpl.java | 4 ++--
2 files changed, 7 insertions(+), 7 deletions(-)
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/consumer/ConsumerImpl.java
b/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/consumer/ConsumerImpl.java
index 0301f0915..d2987d71b 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/consumer/ConsumerImpl.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/consumer/ConsumerImpl.java
@@ -56,13 +56,13 @@ public class ConsumerImpl {
// Other config props
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,
properties.getProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG));
- props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
CloudEventDeserializer.class);
- props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class);
+ props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, new
StringDeserializer());
+ props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, new
CloudEventDeserializer());
props.put(ConsumerConfig.GROUP_ID_CONFIG,
properties.getProperty(ConsumerConfig.GROUP_ID_CONFIG));
- props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
+ props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
this.properties = props;
- this.kafkaConsumer = new KafkaConsumer<String, CloudEvent>(props);
+ this.kafkaConsumer = new KafkaConsumer<>(props);
kafkaConsumerRunner = new KafkaConsumerRunner(this.kafkaConsumer);
executorService = Executors.newFixedThreadPool(10);
topicsSet = new HashSet<>();
@@ -114,7 +114,7 @@ public class ConsumerImpl {
public synchronized void unsubscribe(String topic) {
try {
- // Kafka will unsubscribe *all* topic if calling unsubscribe, so we
+ // Kafka will unsubscribe *all* topic if calling unsubscribe
this.kafkaConsumer.unsubscribe();
topicsSet.remove(topic);
List<String> topics = new ArrayList<>(topicsSet);
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/producer/ProducerImpl.java
b/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/producer/ProducerImpl.java
index a37a30352..0ee990e3b 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/producer/ProducerImpl.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/producer/ProducerImpl.java
@@ -54,8 +54,8 @@ public class ProducerImpl {
properties = new Properties();
properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
props.getProperty(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG));
- properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
StringSerializer.class);
- properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
CloudEventSerializer.class);
+ properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, new
StringSerializer());
+ properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, new
CloudEventSerializer());
this.producer = new KafkaProducer<>(properties);
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]