From ada87dbeee58cfd2fe3453ea1c217cc19036c239 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 20 Jan 2023 14:14:53 -0500 Subject: [PATCH] Pulsar binder schema related changes (#282) * Pulsar binder schema related changes * PR review --- .../src/main/resources/application.yml | 6 +++ .../binder/PulsarMessageChannelBinder.java | 31 +++++------ .../config/PulsarBinderConfiguration.java | 5 +- .../properties/PulsarProducerProperties.java | 6 ++- .../binder/PulsarBinderIntegrationTests.java | 53 ++++++++++++++++++- 5 files changed, 82 insertions(+), 19 deletions(-) diff --git a/spring-pulsar-sample-apps/sample-pulsar-binder/src/main/resources/application.yml b/spring-pulsar-sample-apps/sample-pulsar-binder/src/main/resources/application.yml index 65b1e4ba..106ff7cd 100644 --- a/spring-pulsar-sample-apps/sample-pulsar-binder/src/main/resources/application.yml +++ b/spring-pulsar-sample-apps/sample-pulsar-binder/src/main/resources/application.yml @@ -8,9 +8,15 @@ spring: destination: timeSupplier-out-0 consumer: use-native-decoding: true + timeSupplier-out-0: + producer: + use-native-encoding: true pulsar: bindings: timeLogger-in-0: consumer: subscription-name: my-scst-sub1 schema-type: STRING + timeSupplier-out-0: + producer: + schema-type: STRING \ No newline at end of file diff --git a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarMessageChannelBinder.java b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarMessageChannelBinder.java index 4e16e175..10f9c1c4 100644 --- a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarMessageChannelBinder.java +++ b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarMessageChannelBinder.java @@ -16,6 +16,8 @@ package org.springframework.pulsar.spring.cloud.stream.binder; +import java.util.Objects; + import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.common.schema.SchemaType; @@ -36,6 +38,7 @@ import org.springframework.messaging.MessageHandler; import org.springframework.messaging.support.MessageBuilder; import org.springframework.pulsar.core.PulsarConsumerFactory; import org.springframework.pulsar.core.PulsarTemplate; +import org.springframework.pulsar.core.SchemaResolver; import org.springframework.pulsar.listener.AbstractPulsarMessageListenerContainer; import org.springframework.pulsar.listener.DefaultPulsarMessageListenerContainer; import org.springframework.pulsar.listener.PulsarContainerProperties; @@ -58,19 +61,27 @@ public class PulsarMessageChannelBinder extends private final PulsarConsumerFactory pulsarConsumerFactory; + private final SchemaResolver schemaResolver; + private PulsarExtendedBindingProperties extendedBindingProperties = new PulsarExtendedBindingProperties(); public PulsarMessageChannelBinder(PulsarTopicProvisioner provisioningProvider, - PulsarTemplate pulsarTemplate, PulsarConsumerFactory pulsarConsumerFactory) { + PulsarTemplate pulsarTemplate, PulsarConsumerFactory pulsarConsumerFactory, + SchemaResolver schemaResolver) { super(null, provisioningProvider); this.pulsarTemplate = pulsarTemplate; this.pulsarConsumerFactory = pulsarConsumerFactory; + this.schemaResolver = schemaResolver; } @Override protected MessageHandler createProducerMessageHandler(ProducerDestination destination, ExtendedProducerProperties producerProperties, MessageChannel errorChannel) { - + SchemaType schemaType = producerProperties.getExtension().getSchemaType(); + if (producerProperties.isUseNativeEncoding() && schemaType != null) { + Schema schema = Objects.requireNonNull(this.schemaResolver.getSchema(schemaType, null)); + this.pulsarTemplate.setSchema(schema); + } return message -> { try { PulsarMessageChannelBinder.this.pulsarTemplate.sendAsync(destination.getName(), message.getPayload()); @@ -79,13 +90,11 @@ public class PulsarMessageChannelBinder extends // deal later } }; - } @Override protected MessageProducer createConsumerEndpoint(ConsumerDestination destination, String group, ExtendedConsumerProperties properties) { - PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties.setTopics(new String[] { destination.getName() }); PulsarMessageDrivenChannelAdapter pulsarMessageDrivenChannelAdapter = new PulsarMessageDrivenChannelAdapter(); @@ -94,8 +103,9 @@ public class PulsarMessageChannelBinder extends pulsarMessageDrivenChannelAdapter.send(message); }); SchemaType schemaType = properties.getExtension().getSchemaType(); - if (schemaType != null) { - pulsarContainerProperties.setSchema(toSchema(schemaType)); + if (properties.isUseNativeDecoding() && schemaType != null) { + pulsarContainerProperties + .setSchema(Objects.requireNonNull(this.schemaResolver.getSchema(schemaType, null))); } else { pulsarContainerProperties.setSchema(Schema.BYTES); @@ -107,15 +117,6 @@ public class PulsarMessageChannelBinder extends return pulsarMessageDrivenChannelAdapter; } - // This is just a place-holder impl - public Schema toSchema(SchemaType schemaType) { - return switch (schemaType) { - case STRING -> Schema.STRING; - case INT32 -> Schema.INT32; - default -> Schema.BYTES; - }; - } - @Override public PulsarConsumerProperties getExtendedConsumerProperties(String channelName) { return this.extendedBindingProperties.getExtendedConsumerProperties(channelName); diff --git a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/config/PulsarBinderConfiguration.java b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/config/PulsarBinderConfiguration.java index ce8c6770..508d801d 100644 --- a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/config/PulsarBinderConfiguration.java +++ b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/config/PulsarBinderConfiguration.java @@ -24,6 +24,7 @@ import org.springframework.context.annotation.Configuration; import org.springframework.pulsar.autoconfigure.PulsarProperties; import org.springframework.pulsar.core.PulsarConsumerFactory; import org.springframework.pulsar.core.PulsarTemplate; +import org.springframework.pulsar.core.SchemaResolver; import org.springframework.pulsar.spring.cloud.stream.binder.PulsarMessageChannelBinder; import org.springframework.pulsar.spring.cloud.stream.binder.properties.PulsarExtendedBindingProperties; import org.springframework.pulsar.spring.cloud.stream.binder.provisioning.PulsarTopicProvisioner; @@ -41,9 +42,9 @@ public class PulsarBinderConfiguration { @Bean public PulsarMessageChannelBinder pulsarMessageChannelBinder(PulsarTopicProvisioner pulsarTopicProvisioner, PulsarTemplate pulsarTemplate, PulsarConsumerFactory pulsarConsumerFactory, - PulsarExtendedBindingProperties pulsarExtendedBindingProperties) { + PulsarExtendedBindingProperties pulsarExtendedBindingProperties, SchemaResolver schemaResolver) { PulsarMessageChannelBinder pulsarMessageChannelBinder = new PulsarMessageChannelBinder(pulsarTopicProvisioner, - pulsarTemplate, pulsarConsumerFactory); + pulsarTemplate, pulsarConsumerFactory, schemaResolver); pulsarMessageChannelBinder.setExtendedBindingProperties(pulsarExtendedBindingProperties); return pulsarMessageChannelBinder; } diff --git a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/properties/PulsarProducerProperties.java b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/properties/PulsarProducerProperties.java index 2edfb041..0cdacc80 100644 --- a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/properties/PulsarProducerProperties.java +++ b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/properties/PulsarProducerProperties.java @@ -1,5 +1,5 @@ /* - * Copyright 2022 the original author or authors. + * Copyright 2022-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -18,6 +18,8 @@ package org.springframework.pulsar.spring.cloud.stream.binder.properties; import org.apache.pulsar.common.schema.SchemaType; +import org.springframework.lang.Nullable; + /** * Pulsar producer properties used by the binder. * @@ -25,8 +27,10 @@ import org.apache.pulsar.common.schema.SchemaType; */ public class PulsarProducerProperties { + @Nullable private SchemaType schemaType; + @Nullable public SchemaType getSchemaType() { return this.schemaType; } diff --git a/spring-pulsar-spring-cloud-stream-binder/src/test/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarBinderIntegrationTests.java b/spring-pulsar-spring-cloud-stream-binder/src/test/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarBinderIntegrationTests.java index edec4593..91aa3bbe 100644 --- a/spring-pulsar-spring-cloud-stream-binder/src/test/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarBinderIntegrationTests.java +++ b/spring-pulsar-spring-cloud-stream-binder/src/test/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarBinderIntegrationTests.java @@ -46,7 +46,7 @@ class PulsarBinderIntegrationTests implements PulsarTestContainerSupport { void basicProducerConsumerBindingEndToEnd(CapturedOutput output) { SpringApplication app = new SpringApplication(BasicScenarioConfig.class); app.setWebApplicationType(WebApplicationType.NONE); - try (ConfigurableApplicationContext context = app.run( + try (ConfigurableApplicationContext ignored = app.run( "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), "--spring.cloud.function.definition=textSupplier;textLogger", "--spring.cloud.stream.bindings.textLogger-in-0.destination=textSupplier-out-0", @@ -56,6 +56,39 @@ class PulsarBinderIntegrationTests implements PulsarTestContainerSupport { } } + @Test + void useNativeEncodingDecodingWorkAsExpected(CapturedOutput output) { + SpringApplication app = new SpringApplication(PiStreamConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.cloud.function.definition=piSupplier;piLogger", + "--spring.cloud.stream.bindings.piLogger-in-0.destination=piSupplier-out-0", + "--spring.cloud.stream.bindings.piSupplier-out-0.producer.use-native-encoding=true", + "--spring.cloud.stream.bindings.piLogger-in-0.consumer.use-native-decoding=true", + "--spring.cloud.stream.pulsar.bindings.piLogger-in-0.consumer.schema-type=FLOAT", + "--spring.cloud.stream.pulsar.bindings.piSupplier-out-0.producer.schema-type=FLOAT", + "--spring.cloud.stream.pulsar.bindings.piLogger-in-0.consumer.subscription-name=native-encoding-decoding-sub-1")) { + Awaitility.await().atMost(Duration.ofSeconds(10)) + .until(() -> output.toString().contains("Hello binder: 3.14")); + } + } + + @Test + void basicProducerConsumerBindingEndToEndWithNonTextPayloadType(CapturedOutput output) { + SpringApplication app = new SpringApplication(PiStreamConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.cloud.function.definition=piSupplier;piLogger", + "--spring.cloud.stream.bindings.piSupplier-out-0.destination=pi-stream", + "--spring.cloud.stream.bindings.piLogger-in-0.destination=pi-stream", + "--spring.cloud.stream.pulsar.bindings.piLogger-in-0.consumer.subscription-name=native-encoding-decoding-sub-2")) { + Awaitility.await().atMost(Duration.ofSeconds(10)) + .until(() -> output.toString().contains("Hello binder: 3.14")); + } + } + @EnableAutoConfiguration @SpringBootConfiguration static class BasicScenarioConfig { @@ -74,4 +107,22 @@ class PulsarBinderIntegrationTests implements PulsarTestContainerSupport { } + @EnableAutoConfiguration + @SpringBootConfiguration + static class PiStreamConfig { + + private final Logger logger = LoggerFactory.getLogger(BasicScenarioConfig.class); + + @Bean + public Supplier piSupplier() { + return () -> 3.14f; + } + + @Bean + public Consumer piLogger() { + return f -> this.logger.info("Hello binder: " + f); + } + + } + }