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]