From 3b978984781bd3cd5c4e08bd20de979ed6a789a1 Mon Sep 17 00:00:00 2001 From: wangguowei Date: Tue, 13 Oct 2020 15:01:22 +0800 Subject: [PATCH 1/5] JAVA client support UnAvroBasedSchema --- .../client/api/schema/GenericSchema.java | 13 ++ .../client/impl/schema/AutoConsumeSchema.java | 22 +- .../impl/schema/AvroBasedStructSchema.java | 90 ++++++++ .../pulsar/client/impl/schema/AvroSchema.java | 2 +- .../pulsar/client/impl/schema/JSONSchema.java | 2 +- .../client/impl/schema/ProtobufSchema.java | 2 +- .../client/impl/schema/StructSchema.java | 86 +------- .../AbstractAvroBasedGenericSchema.java | 60 +++++ .../schema/generic/AbstractGenericSchema.java | 54 +++++ .../schema/generic/AvroRecordBuilderImpl.java | 5 +- .../schema/generic/GenericJsonSchema.java | 2 +- .../schema/generic/GenericSchemaImpl.java | 15 +- .../schema/generic/GenericSchemaImplTest.java | 1 + .../schema/generic/GenericSchemaTest.java | 205 ++++++++++++++++++ 14 files changed, 460 insertions(+), 99 deletions(-) create mode 100644 pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AvroBasedStructSchema.java create mode 100644 pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/AbstractAvroBasedGenericSchema.java create mode 100644 pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/AbstractGenericSchema.java create mode 100644 pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/generic/GenericSchemaTest.java 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/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/AvroBasedStructSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AvroBasedStructSchema.java new file mode 100644 index 0000000000000..de86ea3886944 --- /dev/null +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AvroBasedStructSchema.java @@ -0,0 +1,90 @@ +package org.apache.pulsar.client.impl.schema; + +import org.apache.avro.Schema; +import org.apache.avro.reflect.ReflectData; +import org.apache.commons.lang3.StringUtils; +import org.apache.pulsar.client.api.schema.SchemaDefinition; +import org.apache.pulsar.client.impl.schema.generic.GenericAvroSchema; +import org.apache.pulsar.common.schema.SchemaInfo; +import org.apache.pulsar.common.schema.SchemaType; + +import java.lang.reflect.Field; + +import static java.nio.charset.StandardCharsets.UTF_8; + +/** + * Avro Based StructSchema abstract class + */ +public abstract class AvroBasedStructSchema extends StructSchema { + + protected final Schema schema; + + protected AvroBasedStructSchema(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 Schema getAvroSchema() { + return schema; + } + + + protected static Schema createAvroSchema(SchemaDefinition schemaDefinition) { + Class pojo = schemaDefinition.getPojo(); + + if (StringUtils.isNotBlank(schemaDefinition.getJsonDef())) { + return parseAvroSchema(schemaDefinition.getJsonDef()); + } else if (pojo != null) { + ThreadLocal validateDefaults = null; + + try { + Field validateDefaultsField = Schema.class.getDeclaredField("VALIDATE_DEFAULTS"); + validateDefaultsField.setAccessible(true); + validateDefaults = (ThreadLocal) validateDefaultsField.get(null); + } catch (NoSuchFieldException | IllegalAccessException e) { + throw new RuntimeException("Cannot disable validation of default values", e); + } + + final boolean savedValidateDefaults = validateDefaults.get(); + + try { + // Disable validation of default values for compatibility + validateDefaults.set(false); + return extractAvroSchema(schemaDefinition, pojo); + } finally { + validateDefaults.set(savedValidateDefaults); + } + } else { + throw new RuntimeException("Schema definition must specify pojo class or schema json definition"); + } + } + + protected static Schema extractAvroSchema(SchemaDefinition schemaDefinition, Class pojo) { + try { + 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); + } + } + + protected static Schema parseAvroSchema(String schemaJson) { + final Schema.Parser parser = new Schema.Parser(); + parser.setValidateDefaults(false); + return parser.parse(schemaJson); + } + + public static SchemaInfo parseSchemaInfo(SchemaDefinition schemaDefinition, SchemaType schemaType) { + return SchemaInfo.builder() + .schema(createAvroSchema(schemaDefinition).toString().getBytes(UTF_8)) + .properties(schemaDefinition.getProperties()) + .name("") + .type(schemaType).build(); + } + +} diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AvroSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AvroSchema.java index 3e35cb21b287f..15928f7a3e4f3 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AvroSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AvroSchema.java @@ -40,7 +40,7 @@ * An AVRO schema implementation. */ @Slf4j -public class AvroSchema extends StructSchema { +public class AvroSchema extends AvroBasedStructSchema { private static final Logger LOG = LoggerFactory.getLogger(AvroSchema.class); private ClassLoader pojoClassLoader; diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/JSONSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/JSONSchema.java index 45651f3c0e2cc..43abbf6948a00 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/JSONSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/JSONSchema.java @@ -40,7 +40,7 @@ * A schema implementation to deal with json data. */ @Slf4j -public class JSONSchema extends StructSchema { +public class JSONSchema extends AvroBasedStructSchema { // Cannot use org.apache.pulsar.common.util.ObjectMapperFactory.getThreadLocal() because it does not // return shaded version of object mapper private static final ThreadLocal JSON_MAPPER = ThreadLocal.withInitial(() -> { diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ProtobufSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ProtobufSchema.java index 23dde95ed8d2e..24199dc60654a 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ProtobufSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ProtobufSchema.java @@ -45,7 +45,7 @@ /** * A schema implementation to deal with protobuf generated messages. */ -public class ProtobufSchema extends StructSchema { +public class ProtobufSchema extends AvroBasedStructSchema { public static final String PARSING_INFO_PROPERTY = "__PARSING_INFO__"; 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..e5c2480a7ecaa 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,36 +18,26 @@ */ 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.util.concurrent.ExecutionException; +import java.util.concurrent.TimeUnit; + /** * This is a base schema implementation for `Struct` types. * A struct type is used for presenting records (objects) which @@ -62,13 +52,12 @@ public abstract class StructSchema extends AbstractSchema { protected static final Logger LOG = LoggerFactory.getLogger(StructSchema.class); - protected final org.apache.avro.Schema schema; - protected final SchemaInfo schemaInfo; + protected SchemaInfo schemaInfo; protected SchemaReader reader; protected SchemaWriter writer; protected SchemaInfoProvider schemaInfoProvider; - private final LoadingCache> readerCache = CacheBuilder.newBuilder().maximumSize(100000) + LoadingCache> readerCache = CacheBuilder.newBuilder().maximumSize(100000) .expireAfterAccess(30, TimeUnit.MINUTES).build(new CacheLoader>() { @Override public SchemaReader load(BytesSchemaVersion schemaVersion) { @@ -76,19 +65,10 @@ public SchemaReader load(BytesSchemaVersion schemaVersion) { } }); - protected StructSchema(SchemaInfo schemaInfo) { - this.schema = parseAvroSchema(new String(schemaInfo.getSchema(), UTF_8)); + public StructSchema(SchemaInfo schemaInfo){ 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() { - return schema; - } @Override public byte[] encode(T message) { @@ -137,59 +117,7 @@ public SchemaInfo getSchemaInfo() { return this.schemaInfo; } - protected static org.apache.avro.Schema createAvroSchema(SchemaDefinition schemaDefinition) { - Class pojo = schemaDefinition.getPojo(); - - if (StringUtils.isNotBlank(schemaDefinition.getJsonDef())) { - return parseAvroSchema(schemaDefinition.getJsonDef()); - } else if (pojo != null) { - ThreadLocal validateDefaults = null; - - try { - Field validateDefaultsField = Schema.class.getDeclaredField("VALIDATE_DEFAULTS"); - validateDefaultsField.setAccessible(true); - validateDefaults = (ThreadLocal) validateDefaultsField.get(null); - } catch (NoSuchFieldException | IllegalAccessException e) { - throw new RuntimeException("Cannot disable validation of default values", e); - } - - final boolean savedValidateDefaults = validateDefaults.get(); - - try { - // Disable validation of default values for compatibility - validateDefaults.set(false); - return extractAvroSchema(schemaDefinition, pojo); - } finally { - validateDefaults.set(savedValidateDefaults); - } - } else { - throw new RuntimeException("Schema definition must specify pojo class or schema json definition"); - } - } - - protected static Schema extractAvroSchema(SchemaDefinition schemaDefinition, Class pojo) { - try { - 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); - } - } - - protected static org.apache.avro.Schema parseAvroSchema(String schemaJson) { - final Parser parser = new Parser(); - parser.setValidateDefaults(false); - return parser.parse(schemaJson); - } - - public static SchemaInfo parseSchemaInfo(SchemaDefinition schemaDefinition, SchemaType schemaType) { - return SchemaInfo.builder() - .schema(createAvroSchema(schemaDefinition).toString().getBytes(UTF_8)) - .properties(schemaDefinition.getProperties()) - .name("") - .type(schemaType).build(); - } - + @Override public void setSchemaInfoProvider(SchemaInfoProvider schemaInfoProvider) { this.schemaInfoProvider = schemaInfoProvider; } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/AbstractAvroBasedGenericSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/AbstractAvroBasedGenericSchema.java new file mode 100644 index 0000000000000..9e0bec65f6195 --- /dev/null +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/AbstractAvroBasedGenericSchema.java @@ -0,0 +1,60 @@ +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.AvroBasedStructSchema; +import org.apache.pulsar.common.schema.SchemaInfo; + +import java.util.List; +import java.util.stream.Collectors; + +/** + * Abstract AvroBased GenericSchema + */ +abstract class AbstractAvroBasedGenericSchema extends AvroBasedStructSchema implements GenericSchema { + + protected final 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 AbstractAvroBasedGenericSchema(SchemaInfo schemaInfo, boolean useProvidedSchemaAsReaderSchema) { + super(schemaInfo); + this.fields = schema.getFields() + .stream() + .map(f -> new Field(f.name(), f.pos())) + .collect(Collectors.toList()); + this.useProvidedSchemaAsReaderSchema = useProvidedSchemaAsReaderSchema; + } + + @Override + public List getFields() { + return fields; + } + + /** + * Create a generic schema out of a SchemaInfo. + * + * @param schemaInfo schema info + * @return a generic schema instance + */ + public static GenericSchema of(SchemaInfo schemaInfo) { + return of(schemaInfo, true); + } + + public static GenericSchema of(SchemaInfo schemaInfo, + boolean useProvidedSchemaAsReaderSchema) { + switch (schemaInfo.getType()) { + case AVRO: + return new GenericAvroSchema(schemaInfo, useProvidedSchemaAsReaderSchema); + case JSON: + return new GenericJsonSchema(schemaInfo, useProvidedSchemaAsReaderSchema); + default: + throw new UnsupportedOperationException("AvroBased Generic schema is not supported on schema type " + + schemaInfo.getType() + "'"); + } + } + +} 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..6d314e3be2ea3 --- /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.StructSchema; +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 StructSchema 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..76929dc9f7f57 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,17 +22,16 @@ 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}. */ class AvroRecordBuilderImpl implements GenericRecordBuilder { - private final GenericSchemaImpl genericSchema; + private final GenericAvroSchema genericSchema; private final org.apache.avro.generic.GenericRecordBuilder avroRecordBuilder; - AvroRecordBuilderImpl(GenericSchemaImpl genericSchema) { + AvroRecordBuilderImpl(GenericAvroSchema genericSchema) { this.genericSchema = genericSchema; this.avroRecordBuilder = new org.apache.avro.generic.GenericRecordBuilder(genericSchema.getAvroSchema()); 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..a1ae39cdcc42c 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,19 +18,24 @@ */ 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. + * + * warning : + * GenericSchemaImpl will continue to have for backward compatibility and only support AvroBasedGenericSchema ,but qmay deprecate on future , + * we suggest migrate GenericSchemaImpl.of() to .of() method (e.g. GenericJsonSchema 、GenericAvroSchema ) */ -public abstract class GenericSchemaImpl extends StructSchema implements GenericSchema { +@Deprecated +public abstract class GenericSchemaImpl extends AbstractAvroBasedGenericSchema { protected final List fields; // the flag controls whether to use the provided schema as reader schema @@ -40,7 +45,7 @@ public abstract class GenericSchemaImpl extends StructSchema impl protected GenericSchemaImpl(SchemaInfo schemaInfo, boolean useProvidedSchemaAsReaderSchema) { - super(schemaInfo); + super(schemaInfo,useProvidedSchemaAsReaderSchema); this.fields = schema.getFields() .stream() diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/generic/GenericSchemaImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/generic/GenericSchemaImplTest.java index 477e3d808ef4d..df09e05f73267 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/generic/GenericSchemaImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/generic/GenericSchemaImplTest.java @@ -45,6 +45,7 @@ * Unit testing generic schemas. */ @Slf4j +@Deprecated public class GenericSchemaImplTest { @Test 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()); + } + +} From 82561378e8c70caaa8a3981841bbd6bd1b208b53 Mon Sep 17 00:00:00 2001 From: wangguowei Date: Wed, 14 Oct 2020 16:37:35 +0800 Subject: [PATCH 2/5] add AbstractStructSchema for Un-AvroBased StructSchema --- .../impl/schema/AbstractStructSchema.java | 158 +++++++++++++++ .../impl/schema/AvroBasedStructSchema.java | 90 --------- .../pulsar/client/impl/schema/AvroSchema.java | 2 +- .../pulsar/client/impl/schema/JSONSchema.java | 2 +- .../client/impl/schema/KeyValueSchema.java | 4 +- .../client/impl/schema/ProtobufSchema.java | 2 +- .../client/impl/schema/StructSchema.java | 188 ++++++------------ .../AbstractAvroBasedGenericSchema.java | 60 ------ .../schema/generic/AbstractGenericSchema.java | 4 +- .../schema/generic/GenericSchemaImpl.java | 20 +- 10 files changed, 238 insertions(+), 292 deletions(-) create mode 100644 pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AbstractStructSchema.java delete mode 100644 pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AvroBasedStructSchema.java delete mode 100644 pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/AbstractAvroBasedGenericSchema.java 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/AvroBasedStructSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AvroBasedStructSchema.java deleted file mode 100644 index de86ea3886944..0000000000000 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AvroBasedStructSchema.java +++ /dev/null @@ -1,90 +0,0 @@ -package org.apache.pulsar.client.impl.schema; - -import org.apache.avro.Schema; -import org.apache.avro.reflect.ReflectData; -import org.apache.commons.lang3.StringUtils; -import org.apache.pulsar.client.api.schema.SchemaDefinition; -import org.apache.pulsar.client.impl.schema.generic.GenericAvroSchema; -import org.apache.pulsar.common.schema.SchemaInfo; -import org.apache.pulsar.common.schema.SchemaType; - -import java.lang.reflect.Field; - -import static java.nio.charset.StandardCharsets.UTF_8; - -/** - * Avro Based StructSchema abstract class - */ -public abstract class AvroBasedStructSchema extends StructSchema { - - protected final Schema schema; - - protected AvroBasedStructSchema(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 Schema getAvroSchema() { - return schema; - } - - - protected static Schema createAvroSchema(SchemaDefinition schemaDefinition) { - Class pojo = schemaDefinition.getPojo(); - - if (StringUtils.isNotBlank(schemaDefinition.getJsonDef())) { - return parseAvroSchema(schemaDefinition.getJsonDef()); - } else if (pojo != null) { - ThreadLocal validateDefaults = null; - - try { - Field validateDefaultsField = Schema.class.getDeclaredField("VALIDATE_DEFAULTS"); - validateDefaultsField.setAccessible(true); - validateDefaults = (ThreadLocal) validateDefaultsField.get(null); - } catch (NoSuchFieldException | IllegalAccessException e) { - throw new RuntimeException("Cannot disable validation of default values", e); - } - - final boolean savedValidateDefaults = validateDefaults.get(); - - try { - // Disable validation of default values for compatibility - validateDefaults.set(false); - return extractAvroSchema(schemaDefinition, pojo); - } finally { - validateDefaults.set(savedValidateDefaults); - } - } else { - throw new RuntimeException("Schema definition must specify pojo class or schema json definition"); - } - } - - protected static Schema extractAvroSchema(SchemaDefinition schemaDefinition, Class pojo) { - try { - 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); - } - } - - protected static Schema parseAvroSchema(String schemaJson) { - final Schema.Parser parser = new Schema.Parser(); - parser.setValidateDefaults(false); - return parser.parse(schemaJson); - } - - public static SchemaInfo parseSchemaInfo(SchemaDefinition schemaDefinition, SchemaType schemaType) { - return SchemaInfo.builder() - .schema(createAvroSchema(schemaDefinition).toString().getBytes(UTF_8)) - .properties(schemaDefinition.getProperties()) - .name("") - .type(schemaType).build(); - } - -} diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AvroSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AvroSchema.java index 15928f7a3e4f3..3e35cb21b287f 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AvroSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AvroSchema.java @@ -40,7 +40,7 @@ * An AVRO schema implementation. */ @Slf4j -public class AvroSchema extends AvroBasedStructSchema { +public class AvroSchema extends StructSchema { private static final Logger LOG = LoggerFactory.getLogger(AvroSchema.class); private ClassLoader pojoClassLoader; diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/JSONSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/JSONSchema.java index 43abbf6948a00..45651f3c0e2cc 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/JSONSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/JSONSchema.java @@ -40,7 +40,7 @@ * A schema implementation to deal with json data. */ @Slf4j -public class JSONSchema extends AvroBasedStructSchema { +public class JSONSchema extends StructSchema { // Cannot use org.apache.pulsar.common.util.ObjectMapperFactory.getThreadLocal() because it does not // return shaded version of object mapper private static final ThreadLocal JSON_MAPPER = ThreadLocal.withInitial(() -> { 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/ProtobufSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ProtobufSchema.java index 24199dc60654a..23dde95ed8d2e 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ProtobufSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ProtobufSchema.java @@ -45,7 +45,7 @@ /** * A schema implementation to deal with protobuf generated messages. */ -public class ProtobufSchema extends AvroBasedStructSchema { +public class ProtobufSchema extends StructSchema { public static final String PARSING_INFO_PROPERTY = "__PARSING_INFO__"; 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 e5c2480a7ecaa..006f6a97a94d5 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 @@ -1,45 +1,19 @@ -/** - * 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.avro.Schema; +import org.apache.avro.reflect.ReflectData; +import org.apache.commons.lang3.StringUtils; +import org.apache.pulsar.client.api.schema.SchemaDefinition; +import org.apache.pulsar.client.impl.schema.generic.GenericAvroSchema; import org.apache.pulsar.common.schema.SchemaInfo; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; +import org.apache.pulsar.common.schema.SchemaType; + +import java.lang.reflect.Field; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; +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. * @@ -48,118 +22,76 @@ * {@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 SchemaInfo schemaInfo; - protected SchemaReader reader; - protected SchemaWriter writer; - protected SchemaInfoProvider schemaInfoProvider; + protected final Schema schema; - LoadingCache> readerCache = CacheBuilder.newBuilder().maximumSize(100000) - .expireAfterAccess(30, TimeUnit.MINUTES).build(new CacheLoader>() { - @Override - public SchemaReader load(BytesSchemaVersion schemaVersion) { - return loadReader(schemaVersion); - } - }); - - public StructSchema(SchemaInfo schemaInfo){ + 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)); + } } - - @Override - public byte[] encode(T message) { - return writer.write(message); + public Schema getAvroSchema() { + return schema; } - @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()); - } - } + protected static Schema createAvroSchema(SchemaDefinition schemaDefinition) { + Class pojo = schemaDefinition.getPojo(); - @Override - public T decode(ByteBuf byteBuf) { - return reader.read(new ByteBufInputStream(byteBuf)); - } + if (StringUtils.isNotBlank(schemaDefinition.getJsonDef())) { + return parseAvroSchema(schemaDefinition.getJsonDef()); + } else if (pojo != null) { + ThreadLocal validateDefaults = null; - @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()); - } - } + try { + Field validateDefaultsField = Schema.class.getDeclaredField("VALIDATE_DEFAULTS"); + validateDefaultsField.setAccessible(true); + validateDefaults = (ThreadLocal) validateDefaultsField.get(null); + } catch (NoSuchFieldException | IllegalAccessException e) { + throw new RuntimeException("Cannot disable validation of default values", e); + } - @Override - public SchemaInfo getSchemaInfo() { - return this.schemaInfo; - } + final boolean savedValidateDefaults = validateDefaults.get(); - @Override - public void setSchemaInfoProvider(SchemaInfoProvider schemaInfoProvider) { - this.schemaInfoProvider = schemaInfoProvider; + try { + // Disable validation of default values for compatibility + validateDefaults.set(false); + return extractAvroSchema(schemaDefinition, pojo); + } finally { + validateDefaults.set(savedValidateDefaults); + } + } else { + throw new RuntimeException("Schema definition must specify pojo class or schema json definition"); + } } - /** - * 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) { + protected static Schema extractAvroSchema(SchemaDefinition schemaDefinition, Class pojo) { 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() - ); + 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); } } - protected void setWriter(SchemaWriter writer) { - this.writer = writer; - } - - protected void setReader(SchemaReader reader) { - this.reader = reader; + protected static Schema parseAvroSchema(String schemaJson) { + final Schema.Parser parser = new Schema.Parser(); + parser.setValidateDefaults(false); + return parser.parse(schemaJson); } - protected SchemaReader getReader() { - return reader; + public static SchemaInfo parseSchemaInfo(SchemaDefinition schemaDefinition, SchemaType schemaType) { + return SchemaInfo.builder() + .schema(createAvroSchema(schemaDefinition).toString().getBytes(UTF_8)) + .properties(schemaDefinition.getProperties()) + .name("") + .type(schemaType).build(); } } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/AbstractAvroBasedGenericSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/AbstractAvroBasedGenericSchema.java deleted file mode 100644 index 9e0bec65f6195..0000000000000 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/AbstractAvroBasedGenericSchema.java +++ /dev/null @@ -1,60 +0,0 @@ -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.AvroBasedStructSchema; -import org.apache.pulsar.common.schema.SchemaInfo; - -import java.util.List; -import java.util.stream.Collectors; - -/** - * Abstract AvroBased GenericSchema - */ -abstract class AbstractAvroBasedGenericSchema extends AvroBasedStructSchema implements GenericSchema { - - protected final 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 AbstractAvroBasedGenericSchema(SchemaInfo schemaInfo, boolean useProvidedSchemaAsReaderSchema) { - super(schemaInfo); - this.fields = schema.getFields() - .stream() - .map(f -> new Field(f.name(), f.pos())) - .collect(Collectors.toList()); - this.useProvidedSchemaAsReaderSchema = useProvidedSchemaAsReaderSchema; - } - - @Override - public List getFields() { - return fields; - } - - /** - * Create a generic schema out of a SchemaInfo. - * - * @param schemaInfo schema info - * @return a generic schema instance - */ - public static GenericSchema of(SchemaInfo schemaInfo) { - return of(schemaInfo, true); - } - - public static GenericSchema of(SchemaInfo schemaInfo, - boolean useProvidedSchemaAsReaderSchema) { - switch (schemaInfo.getType()) { - case AVRO: - return new GenericAvroSchema(schemaInfo, useProvidedSchemaAsReaderSchema); - case JSON: - return new GenericJsonSchema(schemaInfo, useProvidedSchemaAsReaderSchema); - default: - throw new UnsupportedOperationException("AvroBased Generic schema is not supported on schema type " - + schemaInfo.getType() + "'"); - } - } - -} 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 index 6d314e3be2ea3..00717272f84c3 100644 --- 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 @@ -21,7 +21,7 @@ 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.client.impl.schema.AbstractStructSchema; import org.apache.pulsar.common.schema.SchemaInfo; import java.util.List; @@ -31,7 +31,7 @@ * A minimal abstract generic schema representation for support Un-AvroBasedGenericSchema. * */ -abstract class AbstractGenericSchema extends StructSchema implements GenericSchema { +abstract class AbstractGenericSchema extends AbstractStructSchema implements GenericSchema { protected List fields; // the flag controls whether to use the provided schema as reader schema 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 a1ae39cdcc42c..4f714fe4de28f 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 @@ -28,14 +28,11 @@ import java.util.stream.Collectors; /** - * A generic schema representation. - * + * A generic schema representation for AvroBasedGenericSchema . * warning : - * GenericSchemaImpl will continue to have for backward compatibility and only support AvroBasedGenericSchema ,but qmay deprecate on future , * we suggest migrate GenericSchemaImpl.of() to .of() method (e.g. GenericJsonSchema 、GenericAvroSchema ) */ -@Deprecated -public abstract class GenericSchemaImpl extends AbstractAvroBasedGenericSchema { +public abstract class GenericSchemaImpl extends StructSchema implements GenericSchema { protected final List fields; // the flag controls whether to use the provided schema as reader schema @@ -45,7 +42,7 @@ public abstract class GenericSchemaImpl extends AbstractAvroBasedGenericSchema { protected GenericSchemaImpl(SchemaInfo schemaInfo, boolean useProvidedSchemaAsReaderSchema) { - super(schemaInfo,useProvidedSchemaAsReaderSchema); + super(schemaInfo); this.fields = schema.getFields() .stream() @@ -61,14 +58,23 @@ 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 */ + @Deprecated 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 + */ + @Deprecated public static GenericSchemaImpl of(SchemaInfo schemaInfo, boolean useProvidedSchemaAsReaderSchema) { switch (schemaInfo.getType()) { From 14ff870613b5877ee28c2e1883c860dc3b476c80 Mon Sep 17 00:00:00 2001 From: wangguowei Date: Wed, 14 Oct 2020 16:48:39 +0800 Subject: [PATCH 3/5] add License for StructSchema --- .../client/impl/schema/StructSchema.java | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) 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 006f6a97a94d5..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 @@ -1,3 +1,21 @@ +/** + * 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 org.apache.avro.Schema; From 4074fe7ef0823459b03c01d1cc4022ea8fdfc7df Mon Sep 17 00:00:00 2001 From: wangguowei Date: Fri, 16 Oct 2020 14:46:03 +0800 Subject: [PATCH 4/5] move Deprecated annotation --- .../pulsar/client/impl/schema/generic/GenericSchemaImpl.java | 2 -- .../client/impl/schema/generic/GenericSchemaImplTest.java | 1 - 2 files changed, 3 deletions(-) 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 4f714fe4de28f..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 @@ -62,7 +62,6 @@ public List getFields() { * @param schemaInfo schema info * @return a generic schema instance */ - @Deprecated public static GenericSchemaImpl of(SchemaInfo schemaInfo) { return of(schemaInfo, true); } @@ -74,7 +73,6 @@ public static GenericSchemaImpl of(SchemaInfo schemaInfo) { * @param useProvidedSchemaAsReaderSchema * @return */ - @Deprecated 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/GenericSchemaImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/generic/GenericSchemaImplTest.java index df09e05f73267..477e3d808ef4d 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/generic/GenericSchemaImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/generic/GenericSchemaImplTest.java @@ -45,7 +45,6 @@ * Unit testing generic schemas. */ @Slf4j -@Deprecated public class GenericSchemaImplTest { @Test From 993f22b9d10aecf3746211636a2a83aab8dd11b8 Mon Sep 17 00:00:00 2001 From: wangguowei Date: Mon, 19 Oct 2020 15:41:14 +0800 Subject: [PATCH 5/5] revert GenericAvroSchema to GenericSchemaImpl --- .../client/impl/schema/generic/AvroRecordBuilderImpl.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 76929dc9f7f57..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 @@ -28,10 +28,10 @@ */ class AvroRecordBuilderImpl implements GenericRecordBuilder { - private final GenericAvroSchema genericSchema; + private final GenericSchemaImpl genericSchema; private final org.apache.avro.generic.GenericRecordBuilder avroRecordBuilder; - AvroRecordBuilderImpl(GenericAvroSchema genericSchema) { + AvroRecordBuilderImpl(GenericSchemaImpl genericSchema) { this.genericSchema = genericSchema; this.avroRecordBuilder = new org.apache.avro.generic.GenericRecordBuilder(genericSchema.getAvroSchema());