diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/schema/GenericSchema.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/schema/GenericSchema.java index 95e69393a93ed..4fe293bd171cf 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/schema/GenericSchema.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/schema/GenericSchema.java @@ -20,6 +20,7 @@ import java.util.List; import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.common.schema.SchemaInfo; /** * A schema that serializes and deserializes between {@link GenericRecord} and bytes. @@ -40,4 +41,16 @@ public interface GenericSchema extends Schema { */ GenericRecordBuilder newRecordBuilder(); + + static GenericSchema of(SchemaInfo schemaInfo) { + throw new RuntimeException("GenericSchema interface implementation class must rewrite this method !"); + } + + static GenericSchema of(SchemaInfo schemaInfo, + boolean useProvidedSchemaAsReaderSchema) { + throw new RuntimeException("GenericSchema interface implementation class must rewrite this method !"); + } + + + } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AbstractStructSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AbstractStructSchema.java new file mode 100644 index 0000000000000..7730474707a51 --- /dev/null +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AbstractStructSchema.java @@ -0,0 +1,158 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.client.impl.schema; + +import com.google.common.cache.CacheBuilder; +import com.google.common.cache.CacheLoader; +import com.google.common.cache.LoadingCache; +import io.netty.buffer.ByteBuf; +import io.netty.buffer.ByteBufInputStream; +import org.apache.avro.AvroTypeException; +import org.apache.commons.codec.binary.Hex; +import org.apache.commons.lang3.SerializationException; +import org.apache.pulsar.client.api.SchemaSerializationException; +import org.apache.pulsar.client.api.schema.SchemaInfoProvider; +import org.apache.pulsar.client.api.schema.SchemaReader; +import org.apache.pulsar.client.api.schema.SchemaWriter; +import org.apache.pulsar.common.protocol.schema.BytesSchemaVersion; +import org.apache.pulsar.common.schema.SchemaInfo; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.concurrent.ExecutionException; +import java.util.concurrent.TimeUnit; + +/** + * minimal abstract StructSchema + */ +public abstract class AbstractStructSchema extends AbstractSchema { + + protected static final Logger LOG = LoggerFactory.getLogger(AbstractStructSchema.class); + + protected SchemaInfo schemaInfo; + protected SchemaReader reader; + protected SchemaWriter writer; + protected SchemaInfoProvider schemaInfoProvider; + + LoadingCache> readerCache = CacheBuilder.newBuilder().maximumSize(100000) + .expireAfterAccess(30, TimeUnit.MINUTES).build(new CacheLoader>() { + @Override + public SchemaReader load(BytesSchemaVersion schemaVersion) { + return loadReader(schemaVersion); + } + }); + + public AbstractStructSchema(SchemaInfo schemaInfo){ + this.schemaInfo = schemaInfo; + } + + + @Override + public byte[] encode(T message) { + return writer.write(message); + } + + @Override + public T decode(byte[] bytes) { + return reader.read(bytes); + } + + @Override + public T decode(byte[] bytes, byte[] schemaVersion) { + try { + return schemaVersion == null ? decode(bytes) : + readerCache.get(BytesSchemaVersion.of(schemaVersion)).read(bytes); + } catch (ExecutionException | AvroTypeException e) { + if (e instanceof AvroTypeException) { + throw new SchemaSerializationException(e); + } + LOG.error("Can't get generic schema for topic {} schema version {}", + schemaInfoProvider.getTopicName(), Hex.encodeHexString(schemaVersion), e); + throw new RuntimeException("Can't get generic schema for topic " + schemaInfoProvider.getTopicName()); + } + } + + @Override + public T decode(ByteBuf byteBuf) { + return reader.read(new ByteBufInputStream(byteBuf)); + } + + @Override + public T decode(ByteBuf byteBuf, byte[] schemaVersion) { + try { + return schemaVersion == null ? decode(byteBuf) : + readerCache.get(BytesSchemaVersion.of(schemaVersion)).read(new ByteBufInputStream(byteBuf)); + } catch (ExecutionException e) { + LOG.error("Can't get generic schema for topic {} schema version {}", + schemaInfoProvider.getTopicName(), Hex.encodeHexString(schemaVersion), e); + throw new RuntimeException("Can't get generic schema for topic " + schemaInfoProvider.getTopicName()); + } + } + + @Override + public SchemaInfo getSchemaInfo() { + return this.schemaInfo; + } + + @Override + public void setSchemaInfoProvider(SchemaInfoProvider schemaInfoProvider) { + this.schemaInfoProvider = schemaInfoProvider; + } + + /** + * Load the schema reader for reading messages encoded by the given schema version. + * + * @param schemaVersion the provided schema version + * @return the schema reader for decoding messages encoded by the provided schema version. + */ + protected abstract SchemaReader loadReader(BytesSchemaVersion schemaVersion); + + /** + * TODO: think about how to make this async + */ + protected SchemaInfo getSchemaInfoByVersion(byte[] schemaVersion) { + try { + return schemaInfoProvider.getSchemaByVersion(schemaVersion).get(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new SerializationException( + "Interrupted at fetching schema info for " + SchemaUtils.getStringSchemaVersion(schemaVersion), + e + ); + } catch (ExecutionException e) { + throw new SerializationException( + "Failed at fetching schema info for " + SchemaUtils.getStringSchemaVersion(schemaVersion), + e.getCause() + ); + } + } + + protected void setWriter(SchemaWriter writer) { + this.writer = writer; + } + + protected void setReader(SchemaReader reader) { + this.reader = reader; + } + + protected SchemaReader getReader() { + return reader; + } + +} diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AutoConsumeSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AutoConsumeSchema.java index fb3ac59c31dbe..af902f0a679f6 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AutoConsumeSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AutoConsumeSchema.java @@ -26,10 +26,10 @@ import org.apache.pulsar.client.api.schema.GenericRecord; import org.apache.pulsar.client.api.schema.GenericSchema; import org.apache.pulsar.client.api.schema.SchemaInfoProvider; -import org.apache.pulsar.client.impl.schema.generic.GenericSchemaImpl; +import org.apache.pulsar.client.impl.schema.generic.GenericAvroSchema; +import org.apache.pulsar.client.impl.schema.generic.GenericJsonSchema; import org.apache.pulsar.common.schema.KeyValue; import org.apache.pulsar.common.schema.SchemaInfo; -import org.apache.pulsar.common.schema.SchemaType; import java.util.concurrent.ExecutionException; @@ -147,13 +147,18 @@ public Schema clone() { } private GenericSchema generateSchema(SchemaInfo schemaInfo) { - if (schemaInfo.getType() != SchemaType.AVRO - && schemaInfo.getType() != SchemaType.JSON) { - throw new RuntimeException("Currently auto consume only works for topics with avro or json schemas"); - } // when using `AutoConsumeSchema`, we use the schema associated with the messages as schema reader // to decode the messages. - return GenericSchemaImpl.of(schemaInfo, false /*useProvidedSchemaAsReaderSchema*/); + final boolean useProvidedSchemaAsReaderSchema = false; + switch (schemaInfo.getType()) { + case JSON: + return GenericJsonSchema.of(schemaInfo,useProvidedSchemaAsReaderSchema); + case AVRO: + return GenericAvroSchema.of(schemaInfo,useProvidedSchemaAsReaderSchema); + default: + throw new IllegalArgumentException("Currently auto consume works for type '" + + schemaInfo.getType() + "' is not supported yet"); + } } public static Schema getSchema(SchemaInfo schemaInfo) { @@ -191,8 +196,9 @@ public static Schema getSchema(SchemaInfo schemaInfo) { case LOCAL_DATE_TIME: return LocalDateTimeSchema.of(); case JSON: + return GenericJsonSchema.of(schemaInfo); case AVRO: - return GenericSchemaImpl.of(schemaInfo); + return GenericAvroSchema.of(schemaInfo); case KEY_VALUE: KeyValue kvSchemaInfo = KeyValueSchemaInfo.decodeKeyValueSchemaInfo(schemaInfo); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/KeyValueSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/KeyValueSchema.java index ebc2f0ac6f670..4f7f92173ea00 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/KeyValueSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/KeyValueSchema.java @@ -225,7 +225,7 @@ public CompletableFuture getSchemaByVersion(byte[] schemaVersion) { @Override public CompletableFuture getLatestSchema() { return CompletableFuture.completedFuture( - ((StructSchema) keySchema).schemaInfo); + ((AbstractStructSchema) keySchema).schemaInfo); } @Override @@ -244,7 +244,7 @@ public CompletableFuture getSchemaByVersion(byte[] schemaVersion) { @Override public CompletableFuture getLatestSchema() { return CompletableFuture.completedFuture( - ((StructSchema) valueSchema).schemaInfo); + ((AbstractStructSchema) valueSchema).schemaInfo); } @Override diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/StructSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/StructSchema.java index 6fa344ad52328..986e1866a5305 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/StructSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/StructSchema.java @@ -18,38 +18,20 @@ */ package org.apache.pulsar.client.impl.schema; -import static java.nio.charset.StandardCharsets.UTF_8; - -import java.lang.reflect.Field; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; - -import com.google.common.cache.CacheBuilder; -import com.google.common.cache.CacheLoader; -import com.google.common.cache.LoadingCache; -import io.netty.buffer.ByteBuf; -import io.netty.buffer.ByteBufInputStream; -import org.apache.avro.AvroTypeException; import org.apache.avro.Schema; -import org.apache.avro.Schema.Parser; import org.apache.avro.reflect.ReflectData; -import org.apache.commons.codec.binary.Hex; -import org.apache.commons.lang3.SerializationException; import org.apache.commons.lang3.StringUtils; -import org.apache.pulsar.client.api.SchemaSerializationException; import org.apache.pulsar.client.api.schema.SchemaDefinition; -import org.apache.pulsar.client.api.schema.SchemaInfoProvider; -import org.apache.pulsar.client.api.schema.SchemaReader; -import org.apache.pulsar.client.api.schema.SchemaWriter; import org.apache.pulsar.client.impl.schema.generic.GenericAvroSchema; -import org.apache.pulsar.common.protocol.schema.BytesSchemaVersion; import org.apache.pulsar.common.schema.SchemaInfo; import org.apache.pulsar.common.schema.SchemaType; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; + +import java.lang.reflect.Field; + +import static java.nio.charset.StandardCharsets.UTF_8; /** - * This is a base schema implementation for `Struct` types. + * This is a base schema implementation for Avro Based `Struct` types. * A struct type is used for presenting records (objects) which * have multiple fields. * @@ -58,86 +40,26 @@ * {@link org.apache.pulsar.common.schema.SchemaType#JSON}, * and {@link org.apache.pulsar.common.schema.SchemaType#PROTOBUF}. */ -public abstract class StructSchema extends AbstractSchema { - - protected static final Logger LOG = LoggerFactory.getLogger(StructSchema.class); +public abstract class StructSchema extends AbstractStructSchema { - protected final org.apache.avro.Schema schema; - protected final SchemaInfo schemaInfo; - protected SchemaReader reader; - protected SchemaWriter writer; - protected SchemaInfoProvider schemaInfoProvider; - - private final LoadingCache> readerCache = CacheBuilder.newBuilder().maximumSize(100000) - .expireAfterAccess(30, TimeUnit.MINUTES).build(new CacheLoader>() { - @Override - public SchemaReader load(BytesSchemaVersion schemaVersion) { - return loadReader(schemaVersion); - } - }); + protected final Schema schema; protected StructSchema(SchemaInfo schemaInfo) { + super(schemaInfo); this.schema = parseAvroSchema(new String(schemaInfo.getSchema(), UTF_8)); this.schemaInfo = schemaInfo; - if (schemaInfo.getProperties().containsKey(GenericAvroSchema.OFFSET_PROP)) { this.schema.addProp(GenericAvroSchema.OFFSET_PROP, schemaInfo.getProperties().get(GenericAvroSchema.OFFSET_PROP)); } } - public org.apache.avro.Schema getAvroSchema() { + public Schema getAvroSchema() { return schema; } - @Override - public byte[] encode(T message) { - return writer.write(message); - } - - @Override - public T decode(byte[] bytes) { - return reader.read(bytes); - } - - @Override - public T decode(byte[] bytes, byte[] schemaVersion) { - try { - return schemaVersion == null ? decode(bytes) : - readerCache.get(BytesSchemaVersion.of(schemaVersion)).read(bytes); - } catch (ExecutionException | AvroTypeException e) { - if (e instanceof AvroTypeException) { - throw new SchemaSerializationException(e); - } - LOG.error("Can't get generic schema for topic {} schema version {}", - schemaInfoProvider.getTopicName(), Hex.encodeHexString(schemaVersion), e); - throw new RuntimeException("Can't get generic schema for topic " + schemaInfoProvider.getTopicName()); - } - } - - @Override - public T decode(ByteBuf byteBuf) { - return reader.read(new ByteBufInputStream(byteBuf)); - } - - @Override - public T decode(ByteBuf byteBuf, byte[] schemaVersion) { - try { - return schemaVersion == null ? decode(byteBuf) : - readerCache.get(BytesSchemaVersion.of(schemaVersion)).read(new ByteBufInputStream(byteBuf)); - } catch (ExecutionException e) { - LOG.error("Can't get generic schema for topic {} schema version {}", - schemaInfoProvider.getTopicName(), Hex.encodeHexString(schemaVersion), e); - throw new RuntimeException("Can't get generic schema for topic " + schemaInfoProvider.getTopicName()); - } - } - @Override - public SchemaInfo getSchemaInfo() { - return this.schemaInfo; - } - - protected static org.apache.avro.Schema createAvroSchema(SchemaDefinition schemaDefinition) { + protected static Schema createAvroSchema(SchemaDefinition schemaDefinition) { Class pojo = schemaDefinition.getPojo(); if (StringUtils.isNotBlank(schemaDefinition.getJsonDef())) { @@ -172,12 +94,12 @@ protected static Schema extractAvroSchema(SchemaDefinition schemaDefinition, Cla return parseAvroSchema(pojo.getDeclaredField("SCHEMA$").get(null).toString()); } catch (NoSuchFieldException | IllegalAccessException | IllegalArgumentException ignored) { return schemaDefinition.getAlwaysAllowNull() ? ReflectData.AllowNull.get().getSchema(pojo) - : ReflectData.get().getSchema(pojo); + : ReflectData.get().getSchema(pojo); } } - protected static org.apache.avro.Schema parseAvroSchema(String schemaJson) { - final Parser parser = new Parser(); + protected static Schema parseAvroSchema(String schemaJson) { + final Schema.Parser parser = new Schema.Parser(); parser.setValidateDefaults(false); return parser.parse(schemaJson); } @@ -190,48 +112,4 @@ public static SchemaInfo parseSchemaInfo(SchemaDefinition schemaDefinitio .type(schemaType).build(); } - public void setSchemaInfoProvider(SchemaInfoProvider schemaInfoProvider) { - this.schemaInfoProvider = schemaInfoProvider; - } - - /** - * Load the schema reader for reading messages encoded by the given schema version. - * - * @param schemaVersion the provided schema version - * @return the schema reader for decoding messages encoded by the provided schema version. - */ - protected abstract SchemaReader loadReader(BytesSchemaVersion schemaVersion); - - /** - * TODO: think about how to make this async - */ - protected SchemaInfo getSchemaInfoByVersion(byte[] schemaVersion) { - try { - return schemaInfoProvider.getSchemaByVersion(schemaVersion).get(); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - throw new SerializationException( - "Interrupted at fetching schema info for " + SchemaUtils.getStringSchemaVersion(schemaVersion), - e - ); - } catch (ExecutionException e) { - throw new SerializationException( - "Failed at fetching schema info for " + SchemaUtils.getStringSchemaVersion(schemaVersion), - e.getCause() - ); - } - } - - protected void setWriter(SchemaWriter writer) { - this.writer = writer; - } - - protected void setReader(SchemaReader reader) { - this.reader = reader; - } - - protected SchemaReader getReader() { - return reader; - } - } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/AbstractGenericSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/AbstractGenericSchema.java new file mode 100644 index 0000000000000..00717272f84c3 --- /dev/null +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/AbstractGenericSchema.java @@ -0,0 +1,54 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.client.impl.schema.generic; + +import org.apache.pulsar.client.api.schema.Field; +import org.apache.pulsar.client.api.schema.GenericRecord; +import org.apache.pulsar.client.api.schema.GenericSchema; +import org.apache.pulsar.client.impl.schema.AbstractStructSchema; +import org.apache.pulsar.common.schema.SchemaInfo; + +import java.util.List; + +/** + * + * A minimal abstract generic schema representation for support Un-AvroBasedGenericSchema. + * + */ +abstract class AbstractGenericSchema extends AbstractStructSchema implements GenericSchema { + + protected List fields; + // the flag controls whether to use the provided schema as reader schema + // to decode the messages. In `AUTO_CONSUME` mode, setting this flag to `false` + // allows decoding the messages using the schema associated with the messages. + protected final boolean useProvidedSchemaAsReaderSchema; + + protected AbstractGenericSchema(SchemaInfo schemaInfo, + boolean useProvidedSchemaAsReaderSchema) { + super(schemaInfo); + this.useProvidedSchemaAsReaderSchema = useProvidedSchemaAsReaderSchema; + } + + @Override + public List getFields() { + return fields; + } + + +} diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/AvroRecordBuilderImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/AvroRecordBuilderImpl.java index 7278ab5ac2e33..9a807aceebd80 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/AvroRecordBuilderImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/AvroRecordBuilderImpl.java @@ -22,7 +22,6 @@ import org.apache.pulsar.client.api.schema.GenericRecord; import org.apache.pulsar.client.api.schema.GenericRecordBuilder; -import java.util.List; /** * Builder to build {@link org.apache.pulsar.client.api.schema.GenericRecord}. diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/GenericJsonSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/GenericJsonSchema.java index 2f4db38d6dcfd..8edae25a7227b 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/GenericJsonSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/GenericJsonSchema.java @@ -34,7 +34,7 @@ * A generic json schema. */ @Slf4j -class GenericJsonSchema extends GenericSchemaImpl { +public class GenericJsonSchema extends GenericSchemaImpl { public GenericJsonSchema(SchemaInfo schemaInfo) { this(schemaInfo, true); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/GenericSchemaImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/GenericSchemaImpl.java index 7d18e52979a30..426bb7af3296c 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/GenericSchemaImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/GenericSchemaImpl.java @@ -18,17 +18,19 @@ */ package org.apache.pulsar.client.impl.schema.generic; -import java.util.List; -import java.util.stream.Collectors; - import org.apache.pulsar.client.api.schema.Field; import org.apache.pulsar.client.api.schema.GenericRecord; import org.apache.pulsar.client.api.schema.GenericSchema; import org.apache.pulsar.client.impl.schema.StructSchema; import org.apache.pulsar.common.schema.SchemaInfo; +import java.util.List; +import java.util.stream.Collectors; + /** - * A generic schema representation. + * A generic schema representation for AvroBasedGenericSchema . + * warning : + * we suggest migrate GenericSchemaImpl.of() to .of() method (e.g. GenericJsonSchema 、GenericAvroSchema ) */ public abstract class GenericSchemaImpl extends StructSchema implements GenericSchema { @@ -56,7 +58,7 @@ public List getFields() { /** * Create a generic schema out of a SchemaInfo. - * + * warning : we suggest migrate GenericSchemaImpl.of() to .of() method (e.g. GenericJsonSchema 、GenericAvroSchema ) * @param schemaInfo schema info * @return a generic schema instance */ @@ -64,6 +66,13 @@ public static GenericSchemaImpl of(SchemaInfo schemaInfo) { return of(schemaInfo, true); } + /** + * warning : + * we suggest migrate GenericSchemaImpl.of() to .of() method (e.g. GenericJsonSchema 、GenericAvroSchema ) + * @param schemaInfo + * @param useProvidedSchemaAsReaderSchema + * @return + */ public static GenericSchemaImpl of(SchemaInfo schemaInfo, boolean useProvidedSchemaAsReaderSchema) { switch (schemaInfo.getType()) { diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/generic/GenericSchemaTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/generic/GenericSchemaTest.java new file mode 100644 index 0000000000000..0910e863f9cf9 --- /dev/null +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/generic/GenericSchemaTest.java @@ -0,0 +1,205 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.client.impl.schema.generic; + +import com.google.common.collect.Lists; +import lombok.extern.slf4j.Slf4j; +import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.api.schema.GenericRecord; +import org.apache.pulsar.client.api.schema.GenericSchema; +import org.apache.pulsar.client.impl.schema.AutoConsumeSchema; +import org.apache.pulsar.client.impl.schema.KeyValueSchema; +import org.apache.pulsar.client.impl.schema.KeyValueSchemaInfo; +import org.apache.pulsar.client.impl.schema.SchemaTestUtils.Bar; +import org.apache.pulsar.client.impl.schema.SchemaTestUtils.Foo; +import org.apache.pulsar.common.schema.KeyValue; +import org.apache.pulsar.common.schema.KeyValueEncodingType; +import org.testng.annotations.Test; + +import java.util.List; +import java.util.concurrent.CompletableFuture; + +import static java.nio.charset.StandardCharsets.UTF_8; +import static org.mockito.Mockito.*; +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertTrue; + +/** + * Unit testing generic schemas. + * this test is duplicated with GenericSchemaImplTest independent of GenericSchemaImpl + */ +@Slf4j +public class GenericSchemaTest { + + @Test + public void testGenericAvroSchema() { + Schema encodeSchema = Schema.AVRO(Foo.class); + GenericSchema decodeSchema = GenericAvroSchema.of(encodeSchema.getSchemaInfo()); + testEncodeAndDecodeGenericRecord(encodeSchema, decodeSchema); + } + + @Test + public void testGenericJsonSchema() { + Schema encodeSchema = Schema.JSON(Foo.class); + GenericSchema decodeSchema = GenericJsonSchema.of(encodeSchema.getSchemaInfo()); + testEncodeAndDecodeGenericRecord(encodeSchema, decodeSchema); + } + + @Test + public void testAutoAvroSchema() { + // configure encode schema + Schema encodeSchema = Schema.AVRO(Foo.class); + + // configure the schema info provider + MultiVersionSchemaInfoProvider multiVersionGenericSchemaProvider = mock(MultiVersionSchemaInfoProvider.class); + when(multiVersionGenericSchemaProvider.getSchemaByVersion(any(byte[].class))) + .thenReturn(CompletableFuture.completedFuture(encodeSchema.getSchemaInfo())); + + // configure decode schema + AutoConsumeSchema decodeSchema = new AutoConsumeSchema(); + decodeSchema.configureSchemaInfo( + "test-topic", "topic", encodeSchema.getSchemaInfo() + ); + decodeSchema.setSchemaInfoProvider(multiVersionGenericSchemaProvider); + + testEncodeAndDecodeGenericRecord(encodeSchema, decodeSchema); + } + + @Test + public void testAutoJsonSchema() { + // configure the schema info provider + MultiVersionSchemaInfoProvider multiVersionSchemaInfoProvider = mock(MultiVersionSchemaInfoProvider.class); + GenericSchema genericAvroSchema = GenericAvroSchema.of(Schema.AVRO(Foo.class).getSchemaInfo()); + when(multiVersionSchemaInfoProvider.getSchemaByVersion(any(byte[].class))) + .thenReturn(CompletableFuture.completedFuture(genericAvroSchema.getSchemaInfo())); + + // configure encode schema + Schema encodeSchema = Schema.JSON(Foo.class); + + // configure decode schema + AutoConsumeSchema decodeSchema = new AutoConsumeSchema(); + decodeSchema.configureSchemaInfo("test-topic", "topic", encodeSchema.getSchemaInfo()); + decodeSchema.setSchemaInfoProvider(multiVersionSchemaInfoProvider); + + testEncodeAndDecodeGenericRecord(encodeSchema, decodeSchema); + } + + private void testEncodeAndDecodeGenericRecord(Schema encodeSchema, + Schema decodeSchema) { + int numRecords = 10; + for (int i = 0; i < numRecords; i++) { + Foo foo = newFoo(i); + byte[] data = encodeSchema.encode(foo); + + log.info("Decoding : {}", new String(data, UTF_8)); + + GenericRecord record; + if (decodeSchema instanceof AutoConsumeSchema) { + record = decodeSchema.decode(data, new byte[0]); + } else { + record = decodeSchema.decode(data); + } + verifyFooRecord(record, i); + } + } + + @Test + public void testKeyValueSchema() { + // configure the schema info provider + MultiVersionSchemaInfoProvider multiVersionSchemaInfoProvider = mock(MultiVersionSchemaInfoProvider.class); + GenericSchema genericAvroSchema = GenericAvroSchema.of(Schema.AVRO(Foo.class).getSchemaInfo()); + when(multiVersionSchemaInfoProvider.getSchemaByVersion(any(byte[].class))) + .thenReturn(CompletableFuture.completedFuture( + KeyValueSchemaInfo.encodeKeyValueSchemaInfo( + genericAvroSchema, + genericAvroSchema, + KeyValueEncodingType.INLINE + ) + )); + + List> encodeSchemas = Lists.newArrayList( + Schema.JSON(Foo.class), + Schema.AVRO(Foo.class) + ); + + for (Schema keySchema : encodeSchemas) { + for (Schema valueSchema : encodeSchemas) { + // configure encode schema + Schema> kvSchema = KeyValueSchema.of( + keySchema, valueSchema + ); + + // configure decode schema + Schema> decodeSchema = KeyValueSchema.of( + Schema.AUTO_CONSUME(), Schema.AUTO_CONSUME() + ); + decodeSchema.configureSchemaInfo( + "test-topic", "topic",kvSchema.getSchemaInfo() + ); + decodeSchema.setSchemaInfoProvider(multiVersionSchemaInfoProvider); + + testEncodeAndDecodeKeyValues(kvSchema, decodeSchema); + } + } + + } + + private void testEncodeAndDecodeKeyValues(Schema> encodeSchema, + Schema> decodeSchema) { + int numRecords = 10; + for (int i = 0; i < numRecords; i++) { + Foo foo = newFoo(i); + byte[] data = encodeSchema.encode(new KeyValue<>(foo, foo)); + + KeyValue kv = decodeSchema.decode(data, new byte[0]); + verifyFooRecord(kv.getKey(), i); + verifyFooRecord(kv.getValue(), i); + } + } + + private static Foo newFoo(int i) { + Foo foo = new Foo(); + foo.setField1("field-1-" + i); + foo.setField2("field-2-" + i); + foo.setField3(i); + Bar bar = new Bar(); + bar.setField1(i % 2 == 0); + foo.setField4(bar); + foo.setFieldUnableNull("fieldUnableNull-1-" + i); + + return foo; + } + + private static void verifyFooRecord(GenericRecord record, int i) { + Object field1 = record.getField("field1"); + assertEquals("field-1-" + i, field1, "Field 1 is " + field1.getClass()); + Object field2 = record.getField("field2"); + assertEquals("field-2-" + i, field2, "Field 2 is " + field2.getClass()); + Object field3 = record.getField("field3"); + assertEquals(i, field3, "Field 3 is " + field3.getClass()); + Object field4 = record.getField("field4"); + assertTrue(field4 instanceof GenericRecord); + GenericRecord field4Record = (GenericRecord) field4; + assertEquals(i % 2 == 0, field4Record.getField("field1")); + Object fieldUnableNull = record.getField("fieldUnableNull"); + assertEquals("fieldUnableNull-1-" + i, fieldUnableNull, + "fieldUnableNull 1 is " + fieldUnableNull.getClass()); + } + +}