From ae059822865be6f6b0e95a69a95ad252d8017f0c Mon Sep 17 00:00:00 2001 From: Chris Bono Date: Wed, 16 Oct 2024 10:24:01 -0500 Subject: [PATCH] Update to Pulsar 3.3.2 (#875) This commit updates the version of Pulsar to 3.3.2. Also, in Pulsar 3.3.2 the schema validation of an outgoing message value happens later than it did previously. This requires the PulsarTemplate to widen the try/catch net so that when this happens the producer is closed properly. The backing change in Pulsar 3.3.2 can be seen here https://github.com/apache/pulsar/commit/f3c177e2243e26a7849feb91dbed9fec4c5723c0#diff-095bc2359e03726e031d8c2f210560c6a3218f48a969b0f75a2482e25cd54744R69 --- gradle/libs.versions.toml | 2 +- .../compose.yaml | 2 +- .../compose.yaml | 2 +- .../sample-pulsar-binder/compose.yaml | 2 +- .../download-connectors.sh | 4 +-- .../sample-pulsar-reader/compose.yaml | 2 +- .../sample-reactive/compose.yaml | 2 +- .../pulsar/core/PulsarTemplate.java | 32 ++++++++----------- .../pulsar/docker/standalone/pulsar-start.sh | 2 +- 9 files changed, 23 insertions(+), 27 deletions(-) diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml index 1d03687c..e1ad02ab 100644 --- a/gradle/libs.versions.toml +++ b/gradle/libs.versions.toml @@ -9,7 +9,7 @@ micrometer = "1.14.0-M3" micrometer-docs-gen = "1.0.4" micrometer-tracing = "1.4.0-M3" protobuf = "3.25.5" -pulsar = "3.3.1" +pulsar = "3.3.2" pulsar-reactive = "0.5.7" reactor = "2024.0.0-M6" spring = "6.2.0-RC1" diff --git a/spring-pulsar-sample-apps/sample-failover-custom-router/compose.yaml b/spring-pulsar-sample-apps/sample-failover-custom-router/compose.yaml index 42168e46..79b952d5 100644 --- a/spring-pulsar-sample-apps/sample-failover-custom-router/compose.yaml +++ b/spring-pulsar-sample-apps/sample-failover-custom-router/compose.yaml @@ -1,6 +1,6 @@ services: pulsar: - image: 'apachepulsar/pulsar:3.3.1' + image: 'apachepulsar/pulsar:3.3.2' ports: - '6650' - '8080' diff --git a/spring-pulsar-sample-apps/sample-imperative-produce-consume/compose.yaml b/spring-pulsar-sample-apps/sample-imperative-produce-consume/compose.yaml index 42168e46..79b952d5 100644 --- a/spring-pulsar-sample-apps/sample-imperative-produce-consume/compose.yaml +++ b/spring-pulsar-sample-apps/sample-imperative-produce-consume/compose.yaml @@ -1,6 +1,6 @@ services: pulsar: - image: 'apachepulsar/pulsar:3.3.1' + image: 'apachepulsar/pulsar:3.3.2' ports: - '6650' - '8080' diff --git a/spring-pulsar-sample-apps/sample-pulsar-binder/compose.yaml b/spring-pulsar-sample-apps/sample-pulsar-binder/compose.yaml index 42168e46..79b952d5 100644 --- a/spring-pulsar-sample-apps/sample-pulsar-binder/compose.yaml +++ b/spring-pulsar-sample-apps/sample-pulsar-binder/compose.yaml @@ -1,6 +1,6 @@ services: pulsar: - image: 'apachepulsar/pulsar:3.3.1' + image: 'apachepulsar/pulsar:3.3.2' ports: - '6650' - '8080' diff --git a/spring-pulsar-sample-apps/sample-pulsar-functions/download-connectors.sh b/spring-pulsar-sample-apps/sample-pulsar-functions/download-connectors.sh index c41ac2d2..d6dabb4b 100755 --- a/spring-pulsar-sample-apps/sample-pulsar-functions/download-connectors.sh +++ b/spring-pulsar-sample-apps/sample-pulsar-functions/download-connectors.sh @@ -2,6 +2,6 @@ mkdir connectors cd connectors -wget https://archive.apache.org/dist/pulsar/pulsar-3.3.1/connectors/pulsar-io-cassandra-3.3.1.nar -wget https://archive.apache.org/dist/pulsar/pulsar-3.3.1/connectors/pulsar-io-rabbitmq-3.3.1.nar +wget https://archive.apache.org/dist/pulsar/pulsar-3.3.2/connectors/pulsar-io-cassandra-3.3.2.nar +wget https://archive.apache.org/dist/pulsar/pulsar-3.3.2/connectors/pulsar-io-rabbitmq-3.3.2.nar cd .. diff --git a/spring-pulsar-sample-apps/sample-pulsar-reader/compose.yaml b/spring-pulsar-sample-apps/sample-pulsar-reader/compose.yaml index 42168e46..79b952d5 100644 --- a/spring-pulsar-sample-apps/sample-pulsar-reader/compose.yaml +++ b/spring-pulsar-sample-apps/sample-pulsar-reader/compose.yaml @@ -1,6 +1,6 @@ services: pulsar: - image: 'apachepulsar/pulsar:3.3.1' + image: 'apachepulsar/pulsar:3.3.2' ports: - '6650' - '8080' diff --git a/spring-pulsar-sample-apps/sample-reactive/compose.yaml b/spring-pulsar-sample-apps/sample-reactive/compose.yaml index 42168e46..79b952d5 100644 --- a/spring-pulsar-sample-apps/sample-reactive/compose.yaml +++ b/spring-pulsar-sample-apps/sample-reactive/compose.yaml @@ -1,6 +1,6 @@ services: pulsar: - image: 'apachepulsar/pulsar:3.3.1' + image: 'apachepulsar/pulsar:3.3.2' ports: - '6650' - '8080' diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarTemplate.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarTemplate.java index 73690665..237f19a1 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarTemplate.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarTemplate.java @@ -30,7 +30,6 @@ import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.Schema; -import org.apache.pulsar.client.api.TypedMessageBuilder; import org.apache.pulsar.client.api.interceptor.ProducerInterceptor; import org.apache.pulsar.client.api.transaction.Transaction; @@ -285,25 +284,19 @@ public class PulsarTemplate PulsarMessageSenderContext senderContext = PulsarMessageSenderContext.newContext(topicName, this.beanName); Observation observation = newObservation(senderContext); + Producer producer = null; try { observation.start(); - Producer producer = prepareProducerForSend(topicName, message, schema, encryptionKeys, - producerCustomizer); - TypedMessageBuilder messageBuilder; - try { - var txn = getTransaction(); - messageBuilder = (txn != null) ? producer.newMessage(txn) : producer.newMessage(); - messageBuilder = messageBuilder.value(message); - if (typedMessageBuilderCustomizer != null) { - typedMessageBuilderCustomizer.customize(messageBuilder); - } - // propagate props to message - senderContext.properties().forEach(messageBuilder::property); - } - catch (RuntimeException ex) { - ProducerUtils.closeProducerAsync(producer, this.logger); - throw ex; + producer = prepareProducerForSend(topicName, message, schema, encryptionKeys, producerCustomizer); + var txn = getTransaction(); + var messageBuilder = (txn != null) ? producer.newMessage(txn) : producer.newMessage(); + messageBuilder = messageBuilder.value(message); + if (typedMessageBuilderCustomizer != null) { + typedMessageBuilderCustomizer.customize(messageBuilder); } + // propagate props to message + senderContext.properties().forEach(messageBuilder::property); + var finalProducer = producer; return messageBuilder.sendAsync().whenComplete((msgId, ex) -> { if (ex == null) { this.logger.trace(() -> "Sent msg to '%s' topic".formatted(topicName)); @@ -314,10 +307,13 @@ public class PulsarTemplate observation.error(ex); observation.stop(); } - ProducerUtils.closeProducerAsync(producer, this.logger); + ProducerUtils.closeProducerAsync(finalProducer, this.logger); }); } catch (RuntimeException ex) { + if (producer != null) { + ProducerUtils.closeProducerAsync(producer, this.logger); + } observation.error(ex); observation.stop(); throw ex; diff --git a/tools/pulsar/docker/standalone/pulsar-start.sh b/tools/pulsar/docker/standalone/pulsar-start.sh index 9ff4b5b1..cd59ab03 100755 --- a/tools/pulsar/docker/standalone/pulsar-start.sh +++ b/tools/pulsar/docker/standalone/pulsar-start.sh @@ -3,5 +3,5 @@ docker run -it -p 6650:6650 -p 8080:8080 \ --mount source=pulsardata,target=/pulsar/data \ --mount source=pulsarconf,target=/pulsar/conf \ - apachepulsar/pulsar:3.3.1 \ + apachepulsar/pulsar:3.3.2 \ bin/pulsar standalone