Pulsar binder schema related changes (#282)

* Pulsar binder schema related changes

* PR review
This commit is contained in:
Soby Chacko
2023-01-20 14:14:53 -05:00
committed by GitHub
parent e36ef5f19c
commit ada87dbeee
5 changed files with 82 additions and 19 deletions

View File

@@ -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

View File

@@ -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<Object> pulsarTemplate, PulsarConsumerFactory<?> pulsarConsumerFactory) {
PulsarTemplate<Object> pulsarTemplate, PulsarConsumerFactory<?> pulsarConsumerFactory,
SchemaResolver schemaResolver) {
super(null, provisioningProvider);
this.pulsarTemplate = pulsarTemplate;
this.pulsarConsumerFactory = pulsarConsumerFactory;
this.schemaResolver = schemaResolver;
}
@Override
protected MessageHandler createProducerMessageHandler(ProducerDestination destination,
ExtendedProducerProperties<PulsarProducerProperties> producerProperties, MessageChannel errorChannel) {
SchemaType schemaType = producerProperties.getExtension().getSchemaType();
if (producerProperties.isUseNativeEncoding() && schemaType != null) {
Schema<Object> 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<PulsarConsumerProperties> 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);

View File

@@ -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<Object> pulsarTemplate, PulsarConsumerFactory<byte[]> pulsarConsumerFactory,
PulsarExtendedBindingProperties pulsarExtendedBindingProperties) {
PulsarExtendedBindingProperties pulsarExtendedBindingProperties, SchemaResolver schemaResolver) {
PulsarMessageChannelBinder pulsarMessageChannelBinder = new PulsarMessageChannelBinder(pulsarTopicProvisioner,
pulsarTemplate, pulsarConsumerFactory);
pulsarTemplate, pulsarConsumerFactory, schemaResolver);
pulsarMessageChannelBinder.setExtendedBindingProperties(pulsarExtendedBindingProperties);
return pulsarMessageChannelBinder;
}

View File

@@ -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;
}

View File

@@ -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<Float> piSupplier() {
return () -> 3.14f;
}
@Bean
public Consumer<Float> piLogger() {
return f -> this.logger.info("Hello binder: " + f);
}
}
}