From 52299cba8391b3b11aff031548c8659628dee1a0 Mon Sep 17 00:00:00 2001 From: Andrey Yegorov <8622884+dlg99@users.noreply.github.com> Date: Sat, 23 Oct 2021 21:34:38 -0700 Subject: [PATCH] Stop OffsetStore when stopping the connector (#12457) Source connectors based on KCA (all debezium ones) don't stop properly on error / don't restart. https://github.com/apache/pulsar/pull/12441 fixes one problem, this PR fixes another: ofsetStore is not closed on connector stop() and producer/consumer aren't closed too, preventing the connector from shutting down. Closing offset store on connector stop. (cherry picked from commit 63454e9b2573b8f7ba6a023402c92a17b033ee56) --- .../io/kafka/connect/AbstractKafkaConnectSource.java | 6 ++++++ .../pulsar/io/kafka/connect/PulsarOffsetBackingStore.java | 8 ++++++++ 2 files changed, 14 insertions(+) diff --git a/pulsar-io/kafka-connect-adaptor/src/main/java/org/apache/pulsar/io/kafka/connect/AbstractKafkaConnectSource.java b/pulsar-io/kafka-connect-adaptor/src/main/java/org/apache/pulsar/io/kafka/connect/AbstractKafkaConnectSource.java index 4df0603f8acfc..3f1edef8ec9dd 100644 --- a/pulsar-io/kafka-connect-adaptor/src/main/java/org/apache/pulsar/io/kafka/connect/AbstractKafkaConnectSource.java +++ b/pulsar-io/kafka-connect-adaptor/src/main/java/org/apache/pulsar/io/kafka/connect/AbstractKafkaConnectSource.java @@ -176,6 +176,12 @@ public synchronized Record read() throws Exception { public void close() { if (sourceTask != null) { sourceTask.stop(); + sourceTask = null; + } + + if (offsetStore != null) { + offsetStore.stop(); + offsetStore = null; } } diff --git a/pulsar-io/kafka-connect-adaptor/src/main/java/org/apache/pulsar/io/kafka/connect/PulsarOffsetBackingStore.java b/pulsar-io/kafka-connect-adaptor/src/main/java/org/apache/pulsar/io/kafka/connect/PulsarOffsetBackingStore.java index b63dd0fb0f0d8..ceeb09f1c3fbc 100644 --- a/pulsar-io/kafka-connect-adaptor/src/main/java/org/apache/pulsar/io/kafka/connect/PulsarOffsetBackingStore.java +++ b/pulsar-io/kafka-connect-adaptor/src/main/java/org/apache/pulsar/io/kafka/connect/PulsarOffsetBackingStore.java @@ -169,12 +169,19 @@ public void start() { @Override public void stop() { + log.info("Stopping PulsarOffsetBackingStore"); if (null != producer) { + try { + producer.flush(); + } catch (PulsarClientException pce) { + log.warn("Failed to flush the producer", pce); + } try { producer.close(); } catch (PulsarClientException e) { log.warn("Failed to close producer", e); } + producer = null; } if (null != reader) { try { @@ -182,6 +189,7 @@ public void stop() { } catch (IOException e) { log.warn("Failed to close reader", e); } + reader = null; } if (null != client) { try {