From 55d90518cc1ea016436ec8250c5d668bdc479ff3 Mon Sep 17 00:00:00 2001 From: Jerry Peng Date: Wed, 9 Jun 2021 17:07:12 -0700 Subject: [PATCH 01/26] Remove the unwanted dependencies in the pulsar function's instance jar --- distribution/server/pom.xml | 6 ++++++ pom.xml | 14 +++++++++----- pulsar-functions/runtime-all/pom.xml | 15 ++++++++++++++- .../functions/instance/JavaInstanceMain.java | 1 + 4 files changed, 30 insertions(+), 6 deletions(-) diff --git a/distribution/server/pom.xml b/distribution/server/pom.xml index e61ea42af6bab..e7c406a43b6dc 100644 --- a/distribution/server/pom.xml +++ b/distribution/server/pom.xml @@ -64,6 +64,12 @@ ${project.version} + + org.apache.zookeeper + zookeeper-prometheus-metrics + ${zookeeper.version} + + ${project.groupId} pulsar-package-bookkeeper-storage diff --git a/pom.xml b/pom.xml index 735e28a7316d6..2a641c9d5dec9 100644 --- a/pom.xml +++ b/pom.xml @@ -1200,13 +1200,17 @@ flexible messaging model and an intuitive client API. com.fasterxml.jackson.core * + + org.apache.zookeeper + * + - - org.apache.zookeeper - zookeeper-prometheus-metrics - ${zookeeper.version} - + + + + + diff --git a/pulsar-functions/runtime-all/pom.xml b/pulsar-functions/runtime-all/pom.xml index 556ecb403282b..119fb3ae59240 100644 --- a/pulsar-functions/runtime-all/pom.xml +++ b/pulsar-functions/runtime-all/pom.xml @@ -30,6 +30,19 @@ .. + + pulsar-functions-runtime-all Pulsar Functions :: Runtime All @@ -48,7 +61,7 @@ ${project.groupId} - pulsar-client-original + pulsar-client-api ${project.version} diff --git a/pulsar-functions/runtime-all/src/main/java/org/apache/pulsar/functions/instance/JavaInstanceMain.java b/pulsar-functions/runtime-all/src/main/java/org/apache/pulsar/functions/instance/JavaInstanceMain.java index bd64bf71b68e1..6852792d61486 100644 --- a/pulsar-functions/runtime-all/src/main/java/org/apache/pulsar/functions/instance/JavaInstanceMain.java +++ b/pulsar-functions/runtime-all/src/main/java/org/apache/pulsar/functions/instance/JavaInstanceMain.java @@ -36,6 +36,7 @@ * This class will create three classloaders: * 1. The root classloader that will share interfaces between the function instance * classloader and user code classloader. This classloader will contain the following dependencies + * - pulsar-io-core * - pulsar-functions-api * - pulsar-client-api * - log4j-slf4j-impl From bbd4fd21e7d5b9d912efd8fdfc1ff96ac925578d Mon Sep 17 00:00:00 2001 From: Jerry Peng Date: Wed, 9 Jun 2021 17:31:51 -0700 Subject: [PATCH 02/26] cleaning up --- pom.xml | 5 ----- 1 file changed, 5 deletions(-) diff --git a/pom.xml b/pom.xml index 2a641c9d5dec9..c795a3defb080 100644 --- a/pom.xml +++ b/pom.xml @@ -1206,11 +1206,6 @@ flexible messaging model and an intuitive client API. - - - - - From ae7a45c8846b8531ee30e0e4c80901b9653a69c7 Mon Sep 17 00:00:00 2001 From: Jerry Peng Date: Wed, 9 Jun 2021 21:31:10 -0700 Subject: [PATCH 03/26] fix additional issues that was introduced in earlier PRs --- .../apache/pulsar/functions/runtime/thread/ThreadRuntime.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/thread/ThreadRuntime.java b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/thread/ThreadRuntime.java index e0a9c6472264d..474410c63728c 100644 --- a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/thread/ThreadRuntime.java +++ b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/thread/ThreadRuntime.java @@ -179,7 +179,6 @@ public void start() throws Exception { String.format("%s-%s", FunctionCommon.getFullyQualifiedName(instanceConfig.getFunctionDetails()), instanceConfig.getInstanceId())); - this.fnThread.setContextClassLoader(functionClassLoader); this.fnThread.setUncaughtExceptionHandler(new Thread.UncaughtExceptionHandler() { @Override public void uncaughtException(Thread t, Throwable e) { From acb92e93e4e7d17b78902be61e9c604b65e2a429 Mon Sep 17 00:00:00 2001 From: penghui Date: Thu, 10 Jun 2021 12:19:16 +0800 Subject: [PATCH 04/26] Fix license check. --- pulsar-sql/presto-distribution/LICENSE | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-sql/presto-distribution/LICENSE b/pulsar-sql/presto-distribution/LICENSE index 7450a5c670cae..787886eaffb82 100644 --- a/pulsar-sql/presto-distribution/LICENSE +++ b/pulsar-sql/presto-distribution/LICENSE @@ -450,7 +450,6 @@ The Apache Software License, Version 2.0 * Apache Zookeeper - zookeeper-3.6.3.jar - zookeeper-jute-3.6.3.jar - - zookeeper-prometheus-metrics-3.6.3.jar * Apache Yetus Audience Annotations - audience-annotations-0.5.0.jar * Swagger From e9a4fc5672ae3408b995dfb8856e4e5f57e4dcc9 Mon Sep 17 00:00:00 2001 From: Jerry Peng Date: Wed, 9 Jun 2021 22:34:57 -0700 Subject: [PATCH 05/26] fixing test --- .../tests/integration/io/TestGenericObjectSink.java | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/tests/docker-images/java-test-functions/src/main/java/org/apache/pulsar/tests/integration/io/TestGenericObjectSink.java b/tests/docker-images/java-test-functions/src/main/java/org/apache/pulsar/tests/integration/io/TestGenericObjectSink.java index d131e5b1c1b98..f585d89fbaef3 100644 --- a/tests/docker-images/java-test-functions/src/main/java/org/apache/pulsar/tests/integration/io/TestGenericObjectSink.java +++ b/tests/docker-images/java-test-functions/src/main/java/org/apache/pulsar/tests/integration/io/TestGenericObjectSink.java @@ -56,10 +56,11 @@ public void write(Record record) { if (record.getSchema().getSchemaInfo().getType() == SchemaType.KEY_VALUE) { // assert that we are able to access the schema (leads to ClassCastException if there is a problem) - KeyValueSchema kvSchema = (KeyValueSchema) record.getSchema(); - log.info("key schema type {}", kvSchema.getKeySchema()); - log.info("value schema type {}", kvSchema.getValueSchema()); - log.info("key encoding {}", kvSchema.getKeyValueEncodingType()); + // TODO need to expose KeyValueSchema was an interface in pulsar-client-api +// KeyValueSchema kvSchema = (KeyValueSchema) record.getSchema(); +// log.info("key schema type {}", kvSchema.getKeySchema()); +// log.info("value schema type {}", kvSchema.getValueSchema()); +// log.info("key encoding {}", kvSchema.getKeyValueEncodingType()); KeyValue keyValue = (KeyValue) record.getValue().getNativeObject(); log.info("kvkey {}", keyValue.getKey()); From e3f66b889d3de105242331f8f3b7851684e0bafc Mon Sep 17 00:00:00 2001 From: Jerry Peng Date: Thu, 10 Jun 2021 06:07:20 -0700 Subject: [PATCH 06/26] uncomment code --- .../tests/integration/io/TestGenericObjectSink.java | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/tests/docker-images/java-test-functions/src/main/java/org/apache/pulsar/tests/integration/io/TestGenericObjectSink.java b/tests/docker-images/java-test-functions/src/main/java/org/apache/pulsar/tests/integration/io/TestGenericObjectSink.java index f585d89fbaef3..d131e5b1c1b98 100644 --- a/tests/docker-images/java-test-functions/src/main/java/org/apache/pulsar/tests/integration/io/TestGenericObjectSink.java +++ b/tests/docker-images/java-test-functions/src/main/java/org/apache/pulsar/tests/integration/io/TestGenericObjectSink.java @@ -56,11 +56,10 @@ public void write(Record record) { if (record.getSchema().getSchemaInfo().getType() == SchemaType.KEY_VALUE) { // assert that we are able to access the schema (leads to ClassCastException if there is a problem) - // TODO need to expose KeyValueSchema was an interface in pulsar-client-api -// KeyValueSchema kvSchema = (KeyValueSchema) record.getSchema(); -// log.info("key schema type {}", kvSchema.getKeySchema()); -// log.info("value schema type {}", kvSchema.getValueSchema()); -// log.info("key encoding {}", kvSchema.getKeyValueEncodingType()); + KeyValueSchema kvSchema = (KeyValueSchema) record.getSchema(); + log.info("key schema type {}", kvSchema.getKeySchema()); + log.info("value schema type {}", kvSchema.getValueSchema()); + log.info("key encoding {}", kvSchema.getKeyValueEncodingType()); KeyValue keyValue = (KeyValue) record.getValue().getNativeObject(); log.info("kvkey {}", keyValue.getKey()); From 1b867765a32d36b8a551004c830908b1c2e44a1d Mon Sep 17 00:00:00 2001 From: Jerry Peng Date: Thu, 10 Jun 2021 06:53:09 -0700 Subject: [PATCH 07/26] Fix deps and add test --- .../instance/JavaInstanceDepsTest.java | 74 +++++++++++++++++++ .../docker-images/java-test-functions/pom.xml | 1 + 2 files changed, 75 insertions(+) create mode 100644 pulsar-functions/runtime-all/src/test/java/org/apache/pulsar/functions/instance/JavaInstanceDepsTest.java diff --git a/pulsar-functions/runtime-all/src/test/java/org/apache/pulsar/functions/instance/JavaInstanceDepsTest.java b/pulsar-functions/runtime-all/src/test/java/org/apache/pulsar/functions/instance/JavaInstanceDepsTest.java new file mode 100644 index 0000000000000..9b24e75487357 --- /dev/null +++ b/pulsar-functions/runtime-all/src/test/java/org/apache/pulsar/functions/instance/JavaInstanceDepsTest.java @@ -0,0 +1,74 @@ +/** + * 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.functions.instance; + +import lombok.extern.slf4j.Slf4j; +import org.testng.Assert; +import org.testng.annotations.Test; + +import java.io.File; +import java.io.IOException; +import java.util.Collections; +import java.util.LinkedList; +import java.util.List; +import java.util.zip.ZipEntry; +import java.util.zip.ZipInputStream; + +@Slf4j +/** + * This test serves to make sure that the correct classes are included in the java-instance.jar + * THAT JAR SHOULD ONLY CONTAIN THE INTERFACES THAT PULSAR FUNCTION'S FRAMEWORK USES TO INTERACT WITH USER CODE + * WHICH INCLUDES CLASSES FROM THE FOLLOWING LIBRARIES + * 1. pulsar-io-core + * 2. pulsar-functions-api + * 3. pulsar-client-api + * 4. slf4j-api + * 5. log4j-slf4j-impl + * 6. log4j-api + * 7. log4j-core + */ +public class JavaInstanceDepsTest { + + @Test + public void testInstanceJarDeps() throws IOException { + File jar = new File("target/java-instance.jar"); + ZipInputStream zip = new ZipInputStream(jar.toURI().toURL().openStream()); + + List notAllowedClasses = new LinkedList<>(); + while(true) { + ZipEntry e = zip.getNextEntry(); + if (e == null) + break; + String name = e.getName(); + if (name.endsWith(".class") && !name.startsWith("META-INF")) { + // The only classes in the java-instance.jar should be org.apache.pulsar, slf4j, and log4j classes + // filter out those classes to see if there are any other classes that should not be allowed + if (!name.startsWith("org/apache/pulsar") + && !name.startsWith("org/slf4j") + && !name.startsWith("org/apache/logging/slf4j") + && !name.startsWith("org/apache/logging/log4j")) { + notAllowedClasses.add(name); + } + } + } + + Assert.assertEquals(notAllowedClasses, Collections.emptyList(), notAllowedClasses.toString()); + } +} diff --git a/tests/docker-images/java-test-functions/pom.xml b/tests/docker-images/java-test-functions/pom.xml index b6e98baef9157..9c8864f75ed4a 100644 --- a/tests/docker-images/java-test-functions/pom.xml +++ b/tests/docker-images/java-test-functions/pom.xml @@ -86,6 +86,7 @@ org.apache.pulsar:pulsar-client-api org.apache.pulsar:pulsar-client-admin-api org.apache.pulsar:pulsar-functions-api-examples + com.google.guava:* From 404e66dfb9cba8bf88f867882e81df191c69e82e Mon Sep 17 00:00:00 2001 From: Jerry Peng Date: Thu, 10 Jun 2021 09:12:48 -0700 Subject: [PATCH 08/26] fix deps --- pulsar-functions/runtime-all/pom.xml | 2 +- tests/docker-images/java-test-functions/pom.xml | 1 + 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/pulsar-functions/runtime-all/pom.xml b/pulsar-functions/runtime-all/pom.xml index 119fb3ae59240..9e0219fced064 100644 --- a/pulsar-functions/runtime-all/pom.xml +++ b/pulsar-functions/runtime-all/pom.xml @@ -105,7 +105,7 @@ make-assembly - package + compile single diff --git a/tests/docker-images/java-test-functions/pom.xml b/tests/docker-images/java-test-functions/pom.xml index 9c8864f75ed4a..be52c31533c84 100644 --- a/tests/docker-images/java-test-functions/pom.xml +++ b/tests/docker-images/java-test-functions/pom.xml @@ -87,6 +87,7 @@ org.apache.pulsar:pulsar-client-admin-api org.apache.pulsar:pulsar-functions-api-examples com.google.guava:* + org.apache.commons:* From 16acc3f359831b7e6dc55180cf7f2cdd5a79df42 Mon Sep 17 00:00:00 2001 From: Jerry Peng Date: Thu, 10 Jun 2021 09:54:27 -0700 Subject: [PATCH 09/26] fix pom --- pulsar-functions/runtime-all/pom.xml | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/pulsar-functions/runtime-all/pom.xml b/pulsar-functions/runtime-all/pom.xml index 9e0219fced064..7a8c94e2be9f3 100644 --- a/pulsar-functions/runtime-all/pom.xml +++ b/pulsar-functions/runtime-all/pom.xml @@ -90,6 +90,23 @@ + + + maven-jar-plugin + + + default-jar + none + + unwanted + unwanted + + + + + org.apache.maven.plugins maven-assembly-plugin From 2ab67578b0edce851c36093ccf8512d0f28974fe Mon Sep 17 00:00:00 2001 From: Jerry Peng Date: Thu, 10 Jun 2021 09:55:39 -0700 Subject: [PATCH 10/26] fixing pom --- pulsar-functions/runtime-all/pom.xml | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/pulsar-functions/runtime-all/pom.xml b/pulsar-functions/runtime-all/pom.xml index 7a8c94e2be9f3..d3eaccf8b1942 100644 --- a/pulsar-functions/runtime-all/pom.xml +++ b/pulsar-functions/runtime-all/pom.xml @@ -98,11 +98,7 @@ default-jar - none - - unwanted - unwanted - + From 1696658d375957b064a71aa8026307c1ed620f91 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Thu, 10 Jun 2021 22:23:45 -0700 Subject: [PATCH 11/26] Made SchemaInfo an interface to avoid issues with KeySchema calling back on the wrong classloader --- .../client/admin/internal/SchemasImpl.java | 15 +- .../internal/DefaultImplementation.java | 7 + .../pulsar/common/schema/SchemaInfo.java | 56 ++------ .../client/impl/schema/BooleanSchema.java | 2 +- .../client/impl/schema/ByteBufSchema.java | 2 +- .../client/impl/schema/ByteBufferSchema.java | 2 +- .../pulsar/client/impl/schema/ByteSchema.java | 2 +- .../client/impl/schema/BytesSchema.java | 2 +- .../pulsar/client/impl/schema/DateSchema.java | 2 +- .../client/impl/schema/DoubleSchema.java | 2 +- .../client/impl/schema/FloatSchema.java | 2 +- .../client/impl/schema/InstantSchema.java | 2 +- .../pulsar/client/impl/schema/IntSchema.java | 2 +- .../pulsar/client/impl/schema/JSONSchema.java | 10 +- .../impl/schema/KeyValueSchemaInfo.java | 11 +- .../client/impl/schema/LocalDateSchema.java | 2 +- .../impl/schema/LocalDateTimeSchema.java | 2 +- .../client/impl/schema/LocalTimeSchema.java | 2 +- .../pulsar/client/impl/schema/LongSchema.java | 2 +- .../impl/schema/ProtobufNativeSchema.java | 9 +- .../client/impl/schema/ProtobufSchema.java | 7 +- .../impl/schema/RecordSchemaBuilderImpl.java | 2 +- .../client/impl/schema/SchemaInfoImpl.java | 132 ++++++++++++++++++ .../client/impl/schema/ShortSchema.java | 2 +- .../client/impl/schema/StringSchema.java | 4 +- .../pulsar/client/impl/schema/TimeSchema.java | 2 +- .../client/impl/schema/TimestampSchema.java | 2 +- .../impl/schema/KeyValueSchemaInfoTest.java | 2 +- .../impl/schema/KeyValueSchemaTest.java | 31 ++-- .../client/impl/schema/SchemaInfoTest.java | 12 +- .../client/impl/schema/StringSchemaTest.java | 4 +- .../admin/internal/data/AuthPoliciesImpl.java | 22 ++- .../protocol/schema/SchemaInfoUtil.java | 43 +++--- 33 files changed, 267 insertions(+), 134 deletions(-) create mode 100644 pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/SchemaInfoImpl.java diff --git a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/SchemasImpl.java b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/SchemasImpl.java index 9a9a4ed0bc69a..d4537ca1e172e 100644 --- a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/SchemasImpl.java +++ b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/SchemasImpl.java @@ -31,6 +31,7 @@ import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.admin.Schemas; import org.apache.pulsar.client.api.Authentication; +import org.apache.pulsar.client.impl.schema.SchemaInfoImpl; import org.apache.pulsar.client.internal.DefaultImplementation; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.protocol.schema.DeleteSchemaResponse; @@ -434,7 +435,7 @@ private WebTarget topicPath(TopicName topic, String... parts) { // the util function converts `GetSchemaResponse` to `SchemaInfo` static SchemaInfo convertGetSchemaResponseToSchemaInfo(TopicName tn, GetSchemaResponse response) { - SchemaInfo info = new SchemaInfo(); + byte[] schema; if (response.getType() == SchemaType.KEY_VALUE) { schema = DefaultImplementation.convertKeyValueDataStringToSchemaInfoSchema( @@ -442,11 +443,13 @@ static SchemaInfo convertGetSchemaResponseToSchemaInfo(TopicName tn, } else { schema = response.getData().getBytes(UTF_8); } - info.setSchema(schema); - info.setType(response.getType()); - info.setProperties(response.getProperties()); - info.setName(tn.getLocalName()); - return info; + + return SchemaInfo.builder() + .schema(schema) + .type(response.getType()) + .properties(response.getProperties()) + .name(tn.getLocalName()) + .build(); } static SchemaInfoWithVersion convertGetSchemaResponseToSchemaInfoWithVersion(TopicName tn, diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/internal/DefaultImplementation.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/internal/DefaultImplementation.java index c0fa4d6743a4a..e98cad5bdd9c4 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/internal/DefaultImplementation.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/internal/DefaultImplementation.java @@ -202,6 +202,13 @@ public static Schema newDoubleSchema() { .newInstance()); } + public static SchemaInfo.Builder newSchemaInfoBuilder() { + return catchExceptions( + () -> (SchemaInfo.Builder) getStaticMethod( + "org.apache.pulsar.client.impl.schema.SchemaInfoImpl", "builder", null) + .invoke(null, null)); + } + public static Schema newDateSchema() { return catchExceptions( () -> (Schema) getStaticMethod( diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/common/schema/SchemaInfo.java b/pulsar-client-api/src/main/java/org/apache/pulsar/common/schema/SchemaInfo.java index f2c5860ee0208..cd36863abf80b 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/common/schema/SchemaInfo.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/common/schema/SchemaInfo.java @@ -18,16 +18,7 @@ */ package org.apache.pulsar.common.schema; -import static java.nio.charset.StandardCharsets.UTF_8; -import java.util.Base64; -import java.util.Collections; import java.util.Map; -import lombok.AllArgsConstructor; -import lombok.Builder; -import lombok.Data; -import lombok.EqualsAndHashCode; -import lombok.NoArgsConstructor; -import lombok.experimental.Accessors; import org.apache.pulsar.client.internal.DefaultImplementation; import org.apache.pulsar.common.classification.InterfaceAudience; import org.apache.pulsar.common.classification.InterfaceStability; @@ -37,55 +28,38 @@ */ @InterfaceAudience.Public @InterfaceStability.Stable -@Data -@AllArgsConstructor -@NoArgsConstructor -@Accessors(chain = true) -@Builder -public class SchemaInfo { +public interface SchemaInfo { - @EqualsAndHashCode.Exclude - private String name; + String getName(); /** * The schema data in AVRO JSON format. */ - private byte[] schema; + byte[] getSchema(); /** * The type of schema (AVRO, JSON, PROTOBUF, etc..). */ - private SchemaType type; + SchemaType getType(); /** * Additional properties of the schema definition (implementation defined). */ - @Builder.Default - private Map properties = Collections.emptyMap(); + Map getProperties(); - public String getSchemaDefinition() { - if (null == schema) { - return ""; - } + String getSchemaDefinition(); - switch (type) { - case AVRO: - case JSON: - case PROTOBUF: - case PROTOBUF_NATIVE: - return new String(schema, UTF_8); - case KEY_VALUE: - KeyValue schemaInfoKeyValue = - DefaultImplementation.decodeKeyValueSchemaInfo(this); - return DefaultImplementation.jsonifyKeyValueSchemaInfo(schemaInfoKeyValue); - default: - return Base64.getEncoder().encodeToString(schema); - } + interface Builder { + Builder name(String name); + Builder schema(byte[] schema); + Builder type(SchemaType type); + Builder properties(Map properties); + SchemaInfo build(); } - @Override - public String toString(){ - return DefaultImplementation.jsonifySchemaInfo(this); + static Builder builder() { + return DefaultImplementation.newSchemaInfoBuilder(); } + Builder toBuilder(); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/BooleanSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/BooleanSchema.java index c66ff4332d16a..3b5296ec00301 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/BooleanSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/BooleanSchema.java @@ -32,7 +32,7 @@ public class BooleanSchema extends AbstractSchema { private static final SchemaInfo SCHEMA_INFO; static { - SCHEMA_INFO = new SchemaInfo() + SCHEMA_INFO = new SchemaInfoImpl() .setName("Boolean") .setType(SchemaType.BOOLEAN) .setSchema(new byte[0]); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ByteBufSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ByteBufSchema.java index 658e3984f9e59..ce68298be2b49 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ByteBufSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ByteBufSchema.java @@ -33,7 +33,7 @@ public class ByteBufSchema extends AbstractSchema { private static final SchemaInfo SCHEMA_INFO; static { - SCHEMA_INFO = new SchemaInfo() + SCHEMA_INFO = new SchemaInfoImpl() .setName("ByteBuf") .setType(SchemaType.BYTES) .setSchema(new byte[0]); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ByteBufferSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ByteBufferSchema.java index c560f0e76aa0b..0ff308fe017c0 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ByteBufferSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ByteBufferSchema.java @@ -34,7 +34,7 @@ public class ByteBufferSchema extends AbstractSchema { private static final SchemaInfo SCHEMA_INFO; static { - SCHEMA_INFO = new SchemaInfo() + SCHEMA_INFO = new SchemaInfoImpl() .setName("ByteBuffer") .setType(SchemaType.BYTES) .setSchema(new byte[0]); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ByteSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ByteSchema.java index 4e4c27e07619c..6d516879bd510 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ByteSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ByteSchema.java @@ -32,7 +32,7 @@ public class ByteSchema extends AbstractSchema { private static final SchemaInfo SCHEMA_INFO; static { - SCHEMA_INFO = new SchemaInfo() + SCHEMA_INFO = new SchemaInfoImpl() .setName("INT8") .setType(SchemaType.INT8) .setSchema(new byte[0]); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/BytesSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/BytesSchema.java index 9c7ec373a2820..98a0e66439d3f 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/BytesSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/BytesSchema.java @@ -31,7 +31,7 @@ public class BytesSchema extends AbstractSchema { private static final SchemaInfo SCHEMA_INFO; static { - SCHEMA_INFO = new SchemaInfo() + SCHEMA_INFO = new SchemaInfoImpl() .setName("Bytes") .setType(SchemaType.BYTES) .setSchema(new byte[0]); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/DateSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/DateSchema.java index 295dae6808117..cbdb91202179e 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/DateSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/DateSchema.java @@ -33,7 +33,7 @@ public class DateSchema extends AbstractSchema { private static final SchemaInfo SCHEMA_INFO; static { - SCHEMA_INFO = new SchemaInfo() + SCHEMA_INFO = new SchemaInfoImpl() .setName("Date") .setType(SchemaType.DATE) .setSchema(new byte[0]); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/DoubleSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/DoubleSchema.java index baa1aacf17d7f..4b269a60e2b4a 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/DoubleSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/DoubleSchema.java @@ -32,7 +32,7 @@ public class DoubleSchema extends AbstractSchema { private static final SchemaInfo SCHEMA_INFO; static { - SCHEMA_INFO = new SchemaInfo() + SCHEMA_INFO = new SchemaInfoImpl() .setName("Double") .setType(SchemaType.DOUBLE) .setSchema(new byte[0]); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/FloatSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/FloatSchema.java index aed905b7123b8..84d40735bc186 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/FloatSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/FloatSchema.java @@ -32,7 +32,7 @@ public class FloatSchema extends AbstractSchema { private static final SchemaInfo SCHEMA_INFO; static { - SCHEMA_INFO = new SchemaInfo() + SCHEMA_INFO = new SchemaInfoImpl() .setName("Float") .setType(SchemaType.FLOAT) .setSchema(new byte[0]); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/InstantSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/InstantSchema.java index 5830ceaf571d4..db33de7de4d9b 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/InstantSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/InstantSchema.java @@ -33,7 +33,7 @@ public class InstantSchema extends AbstractSchema { private static final SchemaInfo SCHEMA_INFO; static { - SCHEMA_INFO = new SchemaInfo() + SCHEMA_INFO = new SchemaInfoImpl() .setName("Instant") .setType(SchemaType.INSTANT) .setSchema(new byte[0]); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/IntSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/IntSchema.java index fc8338e45b84b..dfad28082186a 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/IntSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/IntSchema.java @@ -32,7 +32,7 @@ public class IntSchema extends AbstractSchema { private static final SchemaInfo SCHEMA_INFO; static { - SCHEMA_INFO = new SchemaInfo() + SCHEMA_INFO = new SchemaInfoImpl() .setName("INT32") .setType(SchemaType.INT32) .setSchema(new byte[0]); 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 4e3b87441b731..9fe6aed99c782 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 @@ -74,11 +74,11 @@ public SchemaInfo getBackwardsCompatibleJsonSchemaInfo() { ObjectMapper objectMapper = new ObjectMapper(); JsonSchemaGenerator schemaGen = new JsonSchemaGenerator(objectMapper); JsonSchema jsonBackwardsCompatibleSchema = schemaGen.generateSchema(pojo); - backwardsCompatibleSchemaInfo = new SchemaInfo(); - backwardsCompatibleSchemaInfo.setName(""); - backwardsCompatibleSchemaInfo.setProperties(schemaInfo.getProperties()); - backwardsCompatibleSchemaInfo.setType(SchemaType.JSON); - backwardsCompatibleSchemaInfo.setSchema(objectMapper.writeValueAsBytes(jsonBackwardsCompatibleSchema)); + backwardsCompatibleSchemaInfo = new SchemaInfoImpl() + .setName("") + .setProperties(schemaInfo.getProperties()) + .setType(SchemaType.JSON) + .setSchema(objectMapper.writeValueAsBytes(jsonBackwardsCompatibleSchema)); } catch (JsonProcessingException ex) { throw new RuntimeException(ex); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/KeyValueSchemaInfo.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/KeyValueSchemaInfo.java index 95735261f494b..fb341d5bbb6a8 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/KeyValueSchemaInfo.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/KeyValueSchemaInfo.java @@ -169,11 +169,12 @@ public static SchemaInfo encodeKeyValueSchemaInfo(String schemaName, properties.put(KV_ENCODING_TYPE, String.valueOf(keyValueEncodingType)); // generate the final schema info - return new SchemaInfo() - .setName(schemaName) - .setType(SchemaType.KEY_VALUE) - .setSchema(schemaData) - .setProperties(properties); + return SchemaInfoImpl.builder() + .name(schemaName) + .type(SchemaType.KEY_VALUE) + .schema(schemaData) + .properties(properties) + .build(); } private static void encodeSubSchemaInfoToParentSchemaProperties(SchemaInfo schemaInfo, diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/LocalDateSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/LocalDateSchema.java index add6fd28b5b12..18ef3af2d1d6a 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/LocalDateSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/LocalDateSchema.java @@ -32,7 +32,7 @@ public class LocalDateSchema extends AbstractSchema { private static final SchemaInfo SCHEMA_INFO; static { - SCHEMA_INFO = new SchemaInfo() + SCHEMA_INFO = new SchemaInfoImpl() .setName("LocalDate") .setType(SchemaType.LOCAL_DATE) .setSchema(new byte[0]); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/LocalDateTimeSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/LocalDateTimeSchema.java index aa86a1920c538..05b2787fdd607 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/LocalDateTimeSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/LocalDateTimeSchema.java @@ -37,7 +37,7 @@ public class LocalDateTimeSchema extends AbstractSchema { public static final String DELIMITER = ":"; static { - SCHEMA_INFO = new SchemaInfo() + SCHEMA_INFO = new SchemaInfoImpl() .setName("LocalDateTime") .setType(SchemaType.LOCAL_DATE_TIME) .setSchema(new byte[0]); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/LocalTimeSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/LocalTimeSchema.java index 6e2bf627006e7..e53c620f8e0a2 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/LocalTimeSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/LocalTimeSchema.java @@ -32,7 +32,7 @@ public class LocalTimeSchema extends AbstractSchema { private static final SchemaInfo SCHEMA_INFO; static { - SCHEMA_INFO = new SchemaInfo() + SCHEMA_INFO = new SchemaInfoImpl() .setName("LocalTime") .setType(SchemaType.LOCAL_TIME) .setSchema(new byte[0]); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/LongSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/LongSchema.java index f1491f48069ca..deccaf4ded80e 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/LongSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/LongSchema.java @@ -32,7 +32,7 @@ public class LongSchema extends AbstractSchema { private static final SchemaInfo SCHEMA_INFO; static { - SCHEMA_INFO = new SchemaInfo() + SCHEMA_INFO = new SchemaInfoImpl() .setName("INT64") .setType(SchemaType.INT64) .setSchema(new byte[0]); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ProtobufNativeSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ProtobufNativeSchema.java index 385fc41191bc3..dfc592fa2c820 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ProtobufNativeSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ProtobufNativeSchema.java @@ -68,15 +68,12 @@ private static Descriptors.Descriptor createProtobufNativeSchema(Class po } private ProtobufNativeSchema(SchemaInfo schemaInfo, T protoMessageInstance) { - super(schemaInfo); + super(schemaInfo.toBuilder().properties(new HashMap<>(schemaInfo.getProperties())).build()); setReader(new ProtobufNativeReader<>(protoMessageInstance)); setWriter(new ProtobufNativeWriter<>()); // update properties with protobuf related properties - Map allProperties = new HashMap<>(); - allProperties.putAll(schemaInfo.getProperties()); // set protobuf parsing info - allProperties.put(PARSING_INFO_PROPERTY, getParsingInfo(protoMessageInstance)); - schemaInfo.setProperties(allProperties); + schemaInfo.getProperties().put(PARSING_INFO_PROPERTY, getParsingInfo(protoMessageInstance)); } private String getParsingInfo(T protoMessageInstance) { @@ -124,7 +121,7 @@ public static ProtobufNativeSchema of(SchemaDefinition schemaDefinition) } Descriptors.Descriptor descriptor = createProtobufNativeSchema(schemaDefinition.getPojo()); - SchemaInfo schemaInfo = SchemaInfo.builder() + SchemaInfo schemaInfo = SchemaInfoImpl.builder() .schema(ProtobufNativeSchemaUtils.serialize(descriptor)) .type(SchemaType.PROTOBUF_NATIVE) .name("") 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 f7971eb4f17fc..120dd32fa2049 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 @@ -65,15 +65,12 @@ private static org.apache.avro.Schema createProtobufAvroSchema(Class pojo } private ProtobufSchema(SchemaInfo schemaInfo, T protoMessageInstance) { - super(schemaInfo); + super(schemaInfo.toBuilder().properties(new HashMap<>(schemaInfo.getProperties())).build()); setReader(new ProtobufReader<>(protoMessageInstance)); setWriter(new ProtobufWriter<>()); // update properties with protobuf related properties - Map allProperties = new HashMap<>(); - allProperties.putAll(schemaInfo.getProperties()); // set protobuf parsing info - allProperties.put(PARSING_INFO_PROPERTY, getParsingInfo(protoMessageInstance)); - schemaInfo.setProperties(allProperties); + schemaInfo.getProperties().put(PARSING_INFO_PROPERTY, getParsingInfo(protoMessageInstance)); } private String getParsingInfo(T protoMessageInstance) { diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/RecordSchemaBuilderImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/RecordSchemaBuilderImpl.java index ee9f0cb91f9a8..0fda7d52b0dbe 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/RecordSchemaBuilderImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/RecordSchemaBuilderImpl.java @@ -105,7 +105,7 @@ public SchemaInfo build(SchemaType schemaType) { } baseSchema.setFields(avroFields); - return new SchemaInfo( + return new SchemaInfoImpl( name, baseSchema.toString().getBytes(UTF_8), schemaType, diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/SchemaInfoImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/SchemaInfoImpl.java new file mode 100644 index 0000000000000..6ec4309824c95 --- /dev/null +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/SchemaInfoImpl.java @@ -0,0 +1,132 @@ +/** + * 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 static java.nio.charset.StandardCharsets.UTF_8; +import java.util.Base64; +import java.util.Collections; +import java.util.Map; +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.EqualsAndHashCode; +import lombok.NoArgsConstructor; +import lombok.experimental.Accessors; +import org.apache.pulsar.common.classification.InterfaceAudience; +import org.apache.pulsar.common.classification.InterfaceStability; +import org.apache.pulsar.common.schema.KeyValue; +import org.apache.pulsar.common.schema.SchemaInfo; +import org.apache.pulsar.common.schema.SchemaType; + +/** + * Information about the schema. + */ +@InterfaceAudience.Public +@InterfaceStability.Stable +@Data +@AllArgsConstructor +@NoArgsConstructor +@Accessors(chain = true) +public class SchemaInfoImpl implements SchemaInfo { + + @EqualsAndHashCode.Exclude + private String name; + + /** + * The schema data in AVRO JSON format. + */ + private byte[] schema; + + /** + * The type of schema (AVRO, JSON, PROTOBUF, etc..). + */ + private SchemaType type; + + /** + * Additional properties of the schema definition (implementation defined). + */ + private Map properties = Collections.emptyMap(); + + public static SchemaInfoImplBuilder builder() { + return new SchemaInfoImplBuilder(); + } + + @Override + public Builder toBuilder() { + return builder() + .name(name) + .schema(schema) + .type(type) + .properties(properties); + } + + public String getSchemaDefinition() { + if (null == schema) { + return ""; + } + + switch (type) { + case AVRO: + case JSON: + case PROTOBUF: + case PROTOBUF_NATIVE: + return new String(schema, UTF_8); + case KEY_VALUE: + KeyValue schemaInfoKeyValue = KeyValueSchemaInfo.decodeKeyValueSchemaInfo(this); + return SchemaUtils.jsonifyKeyValueSchemaInfo(schemaInfoKeyValue); + default: + return Base64.getEncoder().encodeToString(schema); + } + } + + @Override + public String toString(){ + return SchemaUtils.jsonifySchemaInfo(this); + } + + public static class SchemaInfoImplBuilder implements SchemaInfo.Builder { + private String name; + private byte[] schema; + private SchemaType type; + private Map properties = Collections.emptyMap(); + + public SchemaInfoImplBuilder name(String name) { + this.name = name; + return this; + } + + public SchemaInfoImplBuilder schema(byte[] schema) { + this.schema = schema; + return this; + } + + public SchemaInfoImplBuilder type(SchemaType type) { + this.type = type; + return this; + } + + public SchemaInfoImplBuilder properties(Map properties) { + this.properties = properties; + return this; + } + + public SchemaInfoImpl build() { + return new SchemaInfoImpl(name, schema, type, properties); + } + } +} diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ShortSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ShortSchema.java index 4014405760848..bbb5ad6752938 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ShortSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ShortSchema.java @@ -32,7 +32,7 @@ public class ShortSchema extends AbstractSchema { private static final SchemaInfo SCHEMA_INFO; static { - SCHEMA_INFO = new SchemaInfo() + SCHEMA_INFO = new SchemaInfoImpl() .setName("INT16") .setType(SchemaType.INT16) .setSchema(new byte[0]); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/StringSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/StringSchema.java index 7e57f6ca6ed86..462fa60d89249 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/StringSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/StringSchema.java @@ -46,7 +46,7 @@ public class StringSchema extends AbstractSchema { // Ensure the ordering of the static initialization CHARSET_KEY = "__charset"; DEFAULT_CHARSET = StandardCharsets.UTF_8; - DEFAULT_SCHEMA_INFO = new SchemaInfo() + DEFAULT_SCHEMA_INFO = new SchemaInfoImpl() .setName("String") .setType(SchemaType.STRING) .setSchema(new byte[0]); @@ -87,7 +87,7 @@ public StringSchema(Charset charset) { this.charset = charset; Map properties = new HashMap<>(); properties.put(CHARSET_KEY, charset.name()); - this.schemaInfo = new SchemaInfo() + this.schemaInfo = new SchemaInfoImpl() .setName(DEFAULT_SCHEMA_INFO.getName()) .setType(SchemaType.STRING) .setSchema(DEFAULT_SCHEMA_INFO.getSchema()) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/TimeSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/TimeSchema.java index d56e4da79873f..ab6e1ad438761 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/TimeSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/TimeSchema.java @@ -33,7 +33,7 @@ public class TimeSchema extends AbstractSchema From 1bb315c2c95540dfdfe32570c1be6ef1cafcb239 Mon Sep 17 00:00:00 2001 From: penghui Date: Sat, 12 Jun 2021 12:08:39 +0800 Subject: [PATCH 25/26] Include pulsar-functions-api-examples --- tests/docker-images/java-test-functions/pom.xml | 1 + 1 file changed, 1 insertion(+) diff --git a/tests/docker-images/java-test-functions/pom.xml b/tests/docker-images/java-test-functions/pom.xml index 1be4cba775d7c..539a135d946d8 100644 --- a/tests/docker-images/java-test-functions/pom.xml +++ b/tests/docker-images/java-test-functions/pom.xml @@ -83,6 +83,7 @@ org.apache.pulsar:pulsar-client + org.apache.pulsar:pulsar-functions-api-examples From d77f4b5971076eadd90aef3ebd22a0df78366f9e Mon Sep 17 00:00:00 2001 From: penghui Date: Sat, 12 Jun 2021 14:30:59 +0800 Subject: [PATCH 26/26] Use KeyValueSchema interface to check if the schema of the message is key value schema. --- .../pulsar/client/impl/schema/AutoProduceBytesSchema.java | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AutoProduceBytesSchema.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AutoProduceBytesSchema.java index 3db955484abdd..8971aab421d9e 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AutoProduceBytesSchema.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AutoProduceBytesSchema.java @@ -21,6 +21,7 @@ import static com.google.common.base.Preconditions.checkState; import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.api.schema.KeyValueSchema; import org.apache.pulsar.common.schema.KeyValueEncodingType; import org.apache.pulsar.common.schema.SchemaInfo; import org.apache.pulsar.common.schema.SchemaType; @@ -74,9 +75,9 @@ public byte[] encode(byte[] message) { if (requireSchemaValidation) { // verify if the message can be decoded by the underlying schema - if (schema instanceof KeyValueSchemaImpl - && ((KeyValueSchemaImpl) schema).getKeyValueEncodingType().equals(KeyValueEncodingType.SEPARATED)) { - ((KeyValueSchemaImpl) schema).getValueSchema().validate(message); + if (schema instanceof KeyValueSchema + && ((KeyValueSchema) schema).getKeyValueEncodingType().equals(KeyValueEncodingType.SEPARATED)) { + ((KeyValueSchema) schema).getValueSchema().validate(message); } else { schema.validate(message); }