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]

Reply via email to