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

markli pushed a commit to branch kafka-bug-fix
in repository https://gitbox.apache.org/repos/asf/incubator-eventmesh.git


The following commit(s) were added to refs/heads/kafka-bug-fix by this push:
     new eb4b70209 Fix the initialization of deserializer
eb4b70209 is described below

commit eb4b7020969bd57131b681e44c79e255b51bb4e1
Author: Mark Li <[email protected]>
AuthorDate: Tue Feb 28 15:56:25 2023 -0800

    Fix the initialization of deserializer
---
 .../eventmesh/connector/kafka/consumer/ConsumerImpl.java     |  5 ++++-
 .../connector/kafka/consumer/KafkaConsumerRunner.java        | 12 ++++++------
 2 files changed, 10 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 eb370c459..1faf57a0d 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
@@ -53,12 +53,14 @@ public class ConsumerImpl {
     private Set<String> topicsSet;
 
     public ConsumerImpl(final Properties properties) {
+        ClassLoader original = Thread.currentThread().getContextClassLoader();
+        Thread.currentThread().setContextClassLoader(null);
         Properties props = new Properties();
 
         // 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.VALUE_DESERIALIZER_CLASS_CONFIG, 
CloudEventDeserializer.class);
         props.put(ConsumerConfig.GROUP_ID_CONFIG, 
properties.getProperty(ConsumerConfig.GROUP_ID_CONFIG));
         props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
 
@@ -67,6 +69,7 @@ public class ConsumerImpl {
         kafkaConsumerRunner = new KafkaConsumerRunner(this.kafkaConsumer);
         executorService = Executors.newFixedThreadPool(10);
         topicsSet = new HashSet<>();
+        Thread.currentThread().setContextClassLoader(original);
     }
 
     public Properties attributes() {
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/consumer/KafkaConsumerRunner.java
 
b/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/consumer/KafkaConsumerRunner.java
index bed94c650..381d1f206 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/consumer/KafkaConsumerRunner.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/consumer/KafkaConsumerRunner.java
@@ -62,6 +62,10 @@ public class KafkaConsumerRunner implements Runnable {
     public void run() {
         try {
             while (!closed.get()) {
+                if (consumer.subscription().isEmpty()) {
+                    // consumer cannot poll if it is subscribe to nothing
+                    continue;
+                }
                 ConsumerRecords<String, CloudEvent> records = 
consumer.poll(Duration.ofMillis(10000));
                 // Handle new records
                 records.forEach(rec -> {
@@ -82,8 +86,7 @@ public class KafkaConsumerRunner implements Runnable {
                                         break;
                                     case ManualAck:
                                         // update offset
-                                        log
-                                            .info("message ack, topic: {}, 
current offset:{}", topicName, rec.offset());
+                                        log.info("message ack, topic: {}, 
current offset:{}", topicName, rec.offset());
                                         break;
                                     default:
                                 }
@@ -113,7 +116,4 @@ public class KafkaConsumerRunner implements Runnable {
         closed.set(true);
         consumer.wakeup();
     }
-}
-
-
-
+}
\ No newline at end of file


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to