This is an automated email from the ASF dual-hosted git repository.
mmerli pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-pulsar.git
The following commit(s) were added to refs/heads/master by this push:
new 2dae33d adding avro schema (#1917)
2dae33d is described below
commit 2dae33d6816fd304913bc31352d086e2fde2b38c
Author: Boyang Jerry Peng <[email protected]>
AuthorDate: Mon Jun 11 10:46:09 2018 -0700
adding avro schema (#1917)
* adding avro schema
* improving implementation
* finishing implementation
* remove unnecessary newlines
* fixing poms
* adding avro schema check
* add missing license header
* Add types to proto definitions
* adding compatibiliy unit tests
* shade avro dependencies
* add shading to pulsar client kafka
---
pom.xml | 1 +
.../apache/pulsar/broker/ServiceConfiguration.java | 3 +-
pulsar-broker-shaded/pom.xml | 29 ++++
pulsar-broker/pom.xml | 6 +
.../apache/pulsar/broker/service/ServerCnx.java | 4 +
.../schema/AvroSchemaCompatibilityCheck.java | 87 ++++++++++++
.../service/schema/SchemaRegistryServiceImpl.java | 9 ++
.../src/main/proto/SchemaRegistryFormat.proto | 2 +
.../schema/AvroSchemaCompatibilityCheckTest.java | 149 +++++++++++++++++++++
.../api/SimpleTypedProducerConsumerTest.java | 125 +++++++++++++++++
pulsar-client-admin-shaded/pom.xml | 29 ++++
.../pulsar-client-kafka-shaded/pom.xml | 29 ++++
pulsar-client-shaded/pom.xml | 29 ++++
pulsar-client/pom.xml | 6 +
.../pulsar/client/impl/schema/AvroSchema.java | 98 ++++++++++++++
.../pulsar/client/schemas/AvroSchemaTest.java | 109 +++++++++++++++
.../org/apache/pulsar/common/api/Commands.java | 4 +
.../apache/pulsar/common/api/proto/PulsarApi.java | 6 +
.../apache/pulsar/common/schema/SchemaType.java | 10 +-
pulsar-common/src/main/proto/PulsarApi.proto | 2 +
20 files changed, 735 insertions(+), 2 deletions(-)
diff --git a/pom.xml b/pom.xml
index c99eb80..734c908 100644
--- a/pom.xml
+++ b/pom.xml
@@ -152,6 +152,7 @@ flexible messaging model and an intuitive client
API.</description>
<kafka-client.version>0.10.2.1</kafka-client.version>
<rabbitmq-client.version>5.1.1</rabbitmq-client.version>
<aws-sdk.version>1.11.297</aws-sdk.version>
+ <avro.version>1.8.2</avro.version>
<!-- test dependencies -->
<disruptor.version>3.4.0</disruptor.version>
diff --git
a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java
b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java
index 77e19f5..212c5c4 100644
---
a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java
+++
b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java
@@ -455,7 +455,8 @@ public class ServiceConfiguration implements
PulsarConfiguration {
private String schemaRegistryStorageClassName =
"org.apache.pulsar.broker.service.schema.BookkeeperSchemaStorageFactory";
private Set<String> schemaRegistryCompatibilityCheckers = Sets.newHashSet(
- "org.apache.pulsar.broker.service.schema.JsonSchemaCompatibilityCheck"
+
"org.apache.pulsar.broker.service.schema.JsonSchemaCompatibilityCheck",
+
"org.apache.pulsar.broker.service.schema.AvroSchemaCompatibilityCheck"
);
/**** --- WebSocket --- ****/
diff --git a/pulsar-broker-shaded/pom.xml b/pulsar-broker-shaded/pom.xml
index 181d822..645ab96 100644
--- a/pulsar-broker-shaded/pom.xml
+++ b/pulsar-broker-shaded/pom.xml
@@ -107,6 +107,14 @@
<include>org.apache.httpcomponents:httpclient</include>
<include>commons-logging:commons-logging</include>
<include>org.apache.httpcomponents:httpcore</include>
+ <include>org.apache.avro:avro</include>
+ <!-- Avro transitive dependencies-->
+ <include>org.codehaus.jackson:jackson-core-asl</include>
+ <include>org.codehaus.jackson:jackson-mapper-asl</include>
+ <include>com.thoughtworks.paranamer:paranamer</include>
+ <include>org.xerial.snappy:snappy-java</include>
+ <include>org.apache.commons:commons-compress</include>
+ <include>org.tukaani:xz</include>
</includes>
</artifactSet>
<filters>
@@ -311,6 +319,27 @@
<pattern>org.apache.http</pattern>
<shadedPattern>org.apache.pulsar.shade.org.apache.http</shadedPattern>
</relocation>
+ <relocation>
+ <pattern>org.apache.avro</pattern>
+
<shadedPattern>org.apache.pulsar.shade.org.apache.avro</shadedPattern>
+ </relocation>
+ <!-- Avro transitive dependencies-->
+ <relocation>
+ <pattern>org.codehaus.jackson</pattern>
+
<shadedPattern>org.apache.pulsar.shade.org.codehaus.jackson</shadedPattern>
+ </relocation>
+ <relocation>
+ <pattern>com.thoughtworks.paranamer</pattern>
+
<shadedPattern>org.apache.pulsar.shade.com.thoughtworks.paranamer</shadedPattern>
+ </relocation>
+ <relocation>
+ <pattern>org.xerial.snappy</pattern>
+
<shadedPattern>org.apache.pulsar.shade.org.xerial.snappy</shadedPattern>
+ </relocation>
+ <relocation>
+ <pattern>org.tukaani</pattern>
+
<shadedPattern>org.apache.pulsar.shade.org.tukaani</shadedPattern>
+ </relocation>
</relocations>
</configuration>
</execution>
diff --git a/pulsar-broker/pom.xml b/pulsar-broker/pom.xml
index dd4118d..0a118a8 100644
--- a/pulsar-broker/pom.xml
+++ b/pulsar-broker/pom.xml
@@ -247,6 +247,12 @@
<artifactId>java-semver</artifactId>
</dependency>
+ <dependency>
+ <groupId>org.apache.avro</groupId>
+ <artifactId>avro</artifactId>
+ <version>${avro.version}</version>
+ </dependency>
+
<!-- aspectJ dependencies -->
<dependency>
<groupId>org.aspectj</groupId>
diff --git
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java
index 490f9a4..fcf4045 100644
---
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java
+++
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java
@@ -682,6 +682,10 @@ public class ServerCnx extends PulsarHandler {
return SchemaType.STRING;
case Json:
return SchemaType.JSON;
+ case Protobuf:
+ return SchemaType.PROTOBUF;
+ case Avro:
+ return SchemaType.AVRO;
default:
return SchemaType.NONE;
}
diff --git
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/AvroSchemaCompatibilityCheck.java
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/AvroSchemaCompatibilityCheck.java
new file mode 100644
index 0000000..5d5a77e
--- /dev/null
+++
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/AvroSchemaCompatibilityCheck.java
@@ -0,0 +1,87 @@
+/**
+ * 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.broker.service.schema;
+
+import org.apache.avro.Schema;
+import org.apache.avro.SchemaValidationException;
+import org.apache.avro.SchemaValidator;
+import org.apache.avro.SchemaValidatorBuilder;
+import org.apache.pulsar.common.schema.SchemaData;
+import org.apache.pulsar.common.schema.SchemaType;
+
+
+import java.util.Arrays;
+
+public class AvroSchemaCompatibilityCheck implements SchemaCompatibilityCheck {
+
+ private final CompatibilityStrategy compatibilityStrategy;
+
+ public AvroSchemaCompatibilityCheck () {
+ this(CompatibilityStrategy.FULL);
+ }
+
+ public AvroSchemaCompatibilityCheck(CompatibilityStrategy
compatibilityStrategy) {
+ this.compatibilityStrategy = compatibilityStrategy;
+ }
+
+ @Override
+ public SchemaType getSchemaType() {
+ return SchemaType.AVRO;
+ }
+
+ @Override
+ public boolean isCompatible(SchemaData from, SchemaData to) {
+
+ Schema.Parser fromParser = new Schema.Parser();
+ Schema fromSchema = fromParser.parse(new String(from.getData()));
+ Schema.Parser toParser = new Schema.Parser();
+ Schema toSchema = toParser.parse(new String(to.getData()));
+
+ SchemaValidator schemaValidator =
createSchemaValidator(this.compatibilityStrategy, true);
+ try {
+ schemaValidator.validate(toSchema, Arrays.asList(fromSchema));
+ } catch (SchemaValidationException e) {
+ return false;
+ }
+ return true;
+ }
+
+ public enum CompatibilityStrategy {
+ BACKWARD,
+ FORWARD,
+ FULL
+ }
+
+ private static SchemaValidator createSchemaValidator(CompatibilityStrategy
compatibilityStrategy,
+ boolean onlyLatestValidator)
{
+ final SchemaValidatorBuilder validatorBuilder = new
SchemaValidatorBuilder();
+ switch (compatibilityStrategy) {
+ case BACKWARD:
+ return
createLatestOrAllValidator(validatorBuilder.canReadStrategy(),
onlyLatestValidator);
+ case FORWARD:
+ return
createLatestOrAllValidator(validatorBuilder.canBeReadStrategy(),
onlyLatestValidator);
+ default:
+ return
createLatestOrAllValidator(validatorBuilder.mutualReadStrategy(),
onlyLatestValidator);
+ }
+ }
+
+ private static SchemaValidator
createLatestOrAllValidator(SchemaValidatorBuilder validatorBuilder, boolean
onlyLatest) {
+ return onlyLatest ? validatorBuilder.validateLatest() :
validatorBuilder.validateAll();
+ }
+}
diff --git
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/SchemaRegistryServiceImpl.java
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/SchemaRegistryServiceImpl.java
index b1a7d2c..8cecad4 100644
---
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/SchemaRegistryServiceImpl.java
+++
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/SchemaRegistryServiceImpl.java
@@ -136,6 +136,7 @@ public class SchemaRegistryServiceImpl implements
SchemaRegistryService {
}
private CompletableFuture<Boolean> checkCompatibilityWithLatest(String
schemaId, SchemaData schema) {
+
return getSchema(schemaId).thenApply(storedSchema ->
(storedSchema == null) ||
compatibilityChecks.getOrDefault(
@@ -154,6 +155,10 @@ public class SchemaRegistryServiceImpl implements
SchemaRegistryService {
return SchemaType.STRING;
case JSON:
return SchemaType.JSON;
+ case PROTOBUF:
+ return SchemaType.PROTOBUF;
+ case AVRO:
+ return SchemaType.AVRO;
default:
return SchemaType.NONE;
}
@@ -167,6 +172,10 @@ public class SchemaRegistryServiceImpl implements
SchemaRegistryService {
return SchemaRegistryFormat.SchemaInfo.SchemaType.STRING;
case JSON:
return SchemaRegistryFormat.SchemaInfo.SchemaType.JSON;
+ case PROTOBUF:
+ return SchemaRegistryFormat.SchemaInfo.SchemaType.PROTOBUF;
+ case AVRO:
+ return SchemaRegistryFormat.SchemaInfo.SchemaType.AVRO;
default:
return SchemaRegistryFormat.SchemaInfo.SchemaType.NONE;
}
diff --git a/pulsar-broker/src/main/proto/SchemaRegistryFormat.proto
b/pulsar-broker/src/main/proto/SchemaRegistryFormat.proto
index 8776ddf..90fa145 100644
--- a/pulsar-broker/src/main/proto/SchemaRegistryFormat.proto
+++ b/pulsar-broker/src/main/proto/SchemaRegistryFormat.proto
@@ -27,6 +27,8 @@ message SchemaInfo {
NONE = 1;
STRING = 2;
JSON = 3;
+ PROTOBUF = 4;
+ AVRO = 5;
}
message KeyValuePair {
required string key = 1;
diff --git
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/schema/AvroSchemaCompatibilityCheckTest.java
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/schema/AvroSchemaCompatibilityCheckTest.java
new file mode 100644
index 0000000..de63f54
--- /dev/null
+++
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/schema/AvroSchemaCompatibilityCheckTest.java
@@ -0,0 +1,149 @@
+/**
+ * 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.broker.service.schema;
+
+import org.apache.pulsar.common.schema.SchemaData;
+import org.apache.pulsar.common.schema.SchemaType;
+import org.testng.Assert;
+import org.testng.annotations.Test;
+
+public class AvroSchemaCompatibilityCheckTest {
+
+ private static final String schemaJson1 =
+
"{\"type\":\"record\",\"name\":\"DefaultTest\",\"namespace\":\"org.apache.pulsar.broker.service.schema"
+
+
".AvroSchemaCompatibilityCheckTest$\",\"fields\":[{\"name\":\"field1\",\"type\":\"string\"}]}";
+ private static final SchemaData schemaData1 = getSchemaData(schemaJson1);
+
+ private static final String schemaJson2 =
+
"{\"type\":\"record\",\"name\":\"DefaultTest\",\"namespace\":\"org.apache.pulsar.broker.service.schema"
+
+
".AvroSchemaCompatibilityCheckTest$\",\"fields\":[{\"name\":\"field1\",\"type\":\"string\"},"
+
+
"{\"name\":\"field2\",\"type\":\"string\",\"default\":\"foo\"}]}";
+ private static final SchemaData schemaData2 = getSchemaData(schemaJson2);
+
+ private static final String schemaJson3 =
+
"{\"type\":\"record\",\"name\":\"DefaultTest\",\"namespace\":\"org" +
+
".apache.pulsar.broker.service.schema.AvroSchemaCompatibilityCheckTest$\"," +
+
"\"fields\":[{\"name\":\"field1\",\"type\":\"string\"},{\"name\":\"field2\",\"type\":\"string\"}]}";
+ private static final SchemaData schemaData3 = getSchemaData(schemaJson3);
+
+ private static final String schemaJson4 =
+
"{\"type\":\"record\",\"name\":\"DefaultTest\",\"namespace\":\"org.apache.pulsar.broker.service.schema"
+
+
".AvroSchemaCompatibilityCheckTest$\",\"fields\":[{\"name\":\"field1_v2\",\"type\":\"string\","
+
+ "\"aliases\":[\"field1\"]}]}";
+ private static final SchemaData schemaData4 = getSchemaData(schemaJson4);
+
+ private static final String schemaJson5 =
+
"{\"type\":\"record\",\"name\":\"DefaultTest\",\"namespace\":\"org.apache.pulsar.broker.service.schema"
+
+
".AvroSchemaCompatibilityCheckTest$\",\"fields\":[{\"name\":\"field1\",\"type\":[\"null\","
+
+ "\"string\"]}]}";
+ private static final SchemaData schemaData5 = getSchemaData(schemaJson5);
+
+ private static final String schemaJson6 =
+
"{\"type\":\"record\",\"name\":\"DefaultTest\",\"namespace\":\"org.apache.pulsar.broker.service.schema"
+
+
".AvroSchemaCompatibilityCheckTest$\",\"fields\":[{\"name\":\"field1\",\"type\":[\"null\","
+
+ "\"string\",\"int\"]}]}";
+ private static final SchemaData schemaData6 = getSchemaData(schemaJson6);
+
+ private static final String schemaJson7 =
+
"{\"type\":\"record\",\"name\":\"DefaultTest\",\"namespace\":\"org.apache.pulsar.broker.service.schema"
+
+
".AvroSchemaCompatibilityCheckTest$\",\"fields\":[{\"name\":\"field1\",\"type\":\"string\"},"
+
+
"{\"name\":\"field2\",\"type\":\"string\",\"default\":\"foo\"},{\"name\":\"field3\","
+
+ "\"type\":\"string\",\"default\":\"bar\"}]}";
+ private static final SchemaData schemaData7 = getSchemaData(schemaJson7);
+
+ /**
+ * make sure new schema is backwards compatible with latest
+ */
+ @Test
+ public void testBackwardCompatibility() {
+
+ AvroSchemaCompatibilityCheck avroSchemaCompatibilityCheck = new
AvroSchemaCompatibilityCheck(
+ AvroSchemaCompatibilityCheck.CompatibilityStrategy.BACKWARD
+ );
+
+ // adding a field with default is backwards compatible
+
Assert.assertTrue(avroSchemaCompatibilityCheck.isCompatible(schemaData1,
schemaData2),
+ "adding a field with default is backwards compatible");
+ // adding a field without default is NOT backwards compatible
+
Assert.assertFalse(avroSchemaCompatibilityCheck.isCompatible(schemaData1,
schemaData3),
+ "adding a field without default is NOT backwards compatible");
+ // Modifying a field name is not backwards compatible
+
Assert.assertFalse(avroSchemaCompatibilityCheck.isCompatible(schemaData1,
schemaData4),
+ "Modifying a field name is not backwards compatible");
+ // evolving field to a union is backwards compatible
+
Assert.assertTrue(avroSchemaCompatibilityCheck.isCompatible(schemaData1,
schemaData5),
+ "evolving field to a union is backwards compatible");
+ // removing a field from a union is NOT backwards compatible
+
Assert.assertFalse(avroSchemaCompatibilityCheck.isCompatible(schemaData5,
schemaData1),
+ "removing a field from a union is NOT backwards compatible");
+ // adding a field to a union is backwards compatible
+
Assert.assertTrue(avroSchemaCompatibilityCheck.isCompatible(schemaData5,
schemaData6),
+ "adding a field to a union is backwards compatible");
+ // removing a field a union is NOT backwards compatible
+
Assert.assertFalse(avroSchemaCompatibilityCheck.isCompatible(schemaData6,
schemaData5),
+ "removing a field a union is NOT backwards compatible");
+ }
+
+ /**
+ * Check to make sure the last schema version is forward-compatible with
new schemas
+ */
+ @Test
+ public void testForwardCompatibility() {
+
+ AvroSchemaCompatibilityCheck avroSchemaCompatibilityCheck = new
AvroSchemaCompatibilityCheck(
+ AvroSchemaCompatibilityCheck.CompatibilityStrategy.FORWARD
+ );
+
+
Assert.assertTrue(avroSchemaCompatibilityCheck.isCompatible(schemaData1,
schemaData2),
+ "adding a field is forward compatible");
+
Assert.assertTrue(avroSchemaCompatibilityCheck.isCompatible(schemaData1,
schemaData3),
+ "adding a field is forward compatible");
+
Assert.assertTrue(avroSchemaCompatibilityCheck.isCompatible(schemaData2,
schemaData3),
+ "adding a field is forward compatible");
+
Assert.assertTrue(avroSchemaCompatibilityCheck.isCompatible(schemaData3,
schemaData2),
+ "adding a field is forward compatible");
+
Assert.assertTrue(avroSchemaCompatibilityCheck.isCompatible(schemaData3,
schemaData2),
+ "adding a field is forward compatible");
+
Assert.assertTrue(avroSchemaCompatibilityCheck.isCompatible(schemaData2,
schemaData7),
+ "removing fields is forward compatible");
+
Assert.assertTrue(avroSchemaCompatibilityCheck.isCompatible(schemaData2,
schemaData1),
+ "removing fields with defaults forward compatible");
+ }
+
+ /**
+ * Make sure the new schema is forward- and backward-compatible from the
latest to newest and from the newest to latest.
+ */
+ @Test
+ public void testFullCompatibility() {
+ AvroSchemaCompatibilityCheck avroSchemaCompatibilityCheck = new
AvroSchemaCompatibilityCheck(
+ AvroSchemaCompatibilityCheck.CompatibilityStrategy.FULL
+ );
+
Assert.assertTrue(avroSchemaCompatibilityCheck.isCompatible(schemaData1,
schemaData2),
+ "adding a field with default fully compatible");
+
Assert.assertFalse(avroSchemaCompatibilityCheck.isCompatible(schemaData1,
schemaData3),
+ "adding a field without default is not fully compatible");
+
Assert.assertFalse(avroSchemaCompatibilityCheck.isCompatible(schemaData3,
schemaData1),
+ "adding a field without default is not fully compatible");
+
+ }
+
+ private static SchemaData getSchemaData(String schemaJson) {
+ return
SchemaData.builder().data(schemaJson.getBytes()).type(SchemaType.AVRO).build();
+ }
+}
diff --git
a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleTypedProducerConsumerTest.java
b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleTypedProducerConsumerTest.java
index fd76393..873080f 100644
---
a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleTypedProducerConsumerTest.java
+++
b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleTypedProducerConsumerTest.java
@@ -26,6 +26,7 @@ import java.util.Objects;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import org.apache.pulsar.broker.service.schema.SchemaRegistry;
+import org.apache.pulsar.client.impl.schema.AvroSchema;
import org.apache.pulsar.client.impl.schema.JSONSchema;
import org.apache.pulsar.common.schema.SchemaData;
import org.apache.pulsar.common.schema.SchemaType;
@@ -193,6 +194,130 @@ public class SimpleTypedProducerConsumerTest extends
ProducerConsumerBase {
log.info("-- Exiting {} test --", methodName);
}
+ @Test
+ public void testAvroProducerAndConsumer() throws Exception {
+ log.info("-- Starting {} test --", methodName);
+
+ AvroSchema<AvroEncodedPojo> avroSchema =
+ AvroSchema.of(AvroEncodedPojo.class);
+
+ Consumer<AvroEncodedPojo> consumer = pulsarClient
+ .newConsumer(avroSchema)
+ .topic("persistent://my-property/use/my-ns/my-topic1")
+ .subscriptionName("my-subscriber-name")
+ .subscribe();
+
+ Producer<AvroEncodedPojo> producer = pulsarClient
+ .newProducer(avroSchema)
+ .topic("persistent://my-property/use/my-ns/my-topic1")
+ .create();
+
+ for (int i = 0; i < 10; i++) {
+ String message = "my-message-" + i;
+ producer.send(new AvroEncodedPojo(message));
+ }
+
+ Message<AvroEncodedPojo> msg = null;
+ Set<AvroEncodedPojo> messageSet = Sets.newHashSet();
+ for (int i = 0; i < 10; i++) {
+ msg = consumer.receive(5, TimeUnit.SECONDS);
+ AvroEncodedPojo receivedMessage = msg.getValue();
+ log.debug("Received message: [{}]", receivedMessage);
+ AvroEncodedPojo expectedMessage = new AvroEncodedPojo("my-message-"
+ i);
+ testMessageOrderAndDuplicates(messageSet, receivedMessage,
expectedMessage);
+ }
+ // Acknowledge the consumption of all messages at once
+ consumer.acknowledgeCumulative(msg);
+ consumer.close();
+
+ SchemaRegistry.SchemaAndMetadata storedSchema =
pulsar.getSchemaRegistryService()
+ .getSchema("my-property/my-ns/my-topic1")
+ .get();
+
+ Assert.assertEquals(storedSchema.schema.getData(),
avroSchema.getSchemaInfo().getSchema());
+
+ log.info("-- Exiting {} test --", methodName);
+
+ }
+
+ @Test(expectedExceptions = {PulsarClientException.class})
+ public void testAvroConsumerWithWrongPrestoredSchema() throws Exception {
+ log.info("-- Starting {} test --", methodName);
+
+ byte[] randomSchemaBytes = ("{\n" +
+ " \"type\": \"record\",\n" +
+ " \"namespace\": \"com.example\",\n" +
+ " \"name\": \"FullName\",\n" +
+ " \"fields\": [\n" +
+ " { \"name\": \"first\", \"type\": \"string\" },\n" +
+ " { \"name\": \"last\", \"type\": \"string\" }\n" +
+ " ]\n" +
+ "} ").getBytes();
+
+ pulsar.getSchemaRegistryService()
+ .putSchemaIfAbsent("my-property/my-ns/my-topic1",
+ SchemaData.builder()
+ .type(SchemaType.AVRO)
+ .isDeleted(false)
+ .timestamp(Clock.systemUTC().millis())
+ .user("me")
+ .data(randomSchemaBytes)
+ .props(Collections.emptyMap())
+ .build()
+ ).get();
+
+ Consumer<AvroEncodedPojo> consumer = pulsarClient
+ .newConsumer(AvroSchema.of(AvroEncodedPojo.class))
+ .topic("persistent://my-property/use/my-ns/my-topic1")
+ .subscriptionName("my-subscriber-name")
+ .subscribe();
+
+ log.info("-- Exiting {} test --", methodName);
+ }
+
+ public static class AvroEncodedPojo {
+ private String message;
+
+ public AvroEncodedPojo() {
+ }
+
+ public AvroEncodedPojo(String message) {
+ this.message = message;
+ }
+
+ public String getMessage() {
+ return message;
+ }
+
+ public void setMessage(String message) {
+ this.message = message;
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) {
+ return true;
+ }
+ if (o == null || getClass() != o.getClass()) {
+ return false;
+ }
+ AvroEncodedPojo that = (AvroEncodedPojo) o;
+ return Objects.equals(message, that.message);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(message);
+ }
+
+ @Override
+ public String toString() {
+ return MoreObjects.toStringHelper(this)
+ .add("message", message)
+ .toString();
+ }
+ }
+
public static class JsonEncodedPojo {
private String message;
diff --git a/pulsar-client-admin-shaded/pom.xml
b/pulsar-client-admin-shaded/pom.xml
index 18e1ba9..839091f 100644
--- a/pulsar-client-admin-shaded/pom.xml
+++ b/pulsar-client-admin-shaded/pom.xml
@@ -84,6 +84,14 @@
<include>javax.annotation:*</include>
<include>org.glassfish.hk2*:*</include>
<include>com.fasterxml.jackson.*:*</include>
+ <include>org.apache.avro:avro</include>
+ <!-- Avro transitive dependencies-->
+ <include>org.codehaus.jackson:jackson-core-asl</include>
+ <include>org.codehaus.jackson:jackson-mapper-asl</include>
+ <include>com.thoughtworks.paranamer:paranamer</include>
+ <include>org.xerial.snappy:snappy-java</include>
+ <include>org.apache.commons:commons-compress</include>
+ <include>org.tukaani:xz</include>
</includes>
</artifactSet>
<filters>
@@ -187,6 +195,27 @@
<pattern>org.reactivestreams</pattern>
<shadedPattern>org.apache.pulsar.admin.shade.org.reactivestreams</shadedPattern>
</relocation>
+ <relocation>
+ <pattern>org.apache.avro</pattern>
+
<shadedPattern>org.apache.pulsar.shade.org.apache.avro</shadedPattern>
+ </relocation>
+ <!-- Avro transitive dependencies-->
+ <relocation>
+ <pattern>org.codehaus.jackson</pattern>
+
<shadedPattern>org.apache.pulsar.shade.org.codehaus.jackson</shadedPattern>
+ </relocation>
+ <relocation>
+ <pattern>com.thoughtworks.paranamer</pattern>
+
<shadedPattern>org.apache.pulsar.shade.com.thoughtworks.paranamer</shadedPattern>
+ </relocation>
+ <relocation>
+ <pattern>org.xerial.snappy</pattern>
+
<shadedPattern>org.apache.pulsar.shade.org.xerial.snappy</shadedPattern>
+ </relocation>
+ <relocation>
+ <pattern>org.tukaani</pattern>
+
<shadedPattern>org.apache.pulsar.shade.org.tukaani</shadedPattern>
+ </relocation>
</relocations>
<transformers>
<transformer
implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer"
/>
diff --git a/pulsar-client-kafka-compat/pulsar-client-kafka-shaded/pom.xml
b/pulsar-client-kafka-compat/pulsar-client-kafka-shaded/pom.xml
index e3c6ddb..0bd9453 100644
--- a/pulsar-client-kafka-compat/pulsar-client-kafka-shaded/pom.xml
+++ b/pulsar-client-kafka-compat/pulsar-client-kafka-shaded/pom.xml
@@ -91,6 +91,14 @@
<include>org.apache.httpcomponents:httpclient</include>
<include>commons-logging:commons-logging</include>
<include>org.apache.httpcomponents:httpcore</include>
+ <include>org.apache.avro:avro</include>
+ <!-- Avro transitive dependencies-->
+ <include>org.codehaus.jackson:jackson-core-asl</include>
+ <include>org.codehaus.jackson:jackson-mapper-asl</include>
+ <include>com.thoughtworks.paranamer:paranamer</include>
+ <include>org.xerial.snappy:snappy-java</include>
+ <include>org.apache.commons:commons-compress</include>
+ <include>org.tukaani:xz</include>
</includes>
</artifactSet>
<filters>
@@ -173,6 +181,27 @@
<pattern>org.apache.http</pattern>
<shadedPattern>org.apache.pulsar.shade.org.apache.http</shadedPattern>
</relocation>
+ <relocation>
+ <pattern>org.apache.avro</pattern>
+
<shadedPattern>org.apache.pulsar.shade.org.apache.avro</shadedPattern>
+ </relocation>
+ <!-- Avro transitive dependencies-->
+ <relocation>
+ <pattern>org.codehaus.jackson</pattern>
+
<shadedPattern>org.apache.pulsar.shade.org.codehaus.jackson</shadedPattern>
+ </relocation>
+ <relocation>
+ <pattern>com.thoughtworks.paranamer</pattern>
+
<shadedPattern>org.apache.pulsar.shade.com.thoughtworks.paranamer</shadedPattern>
+ </relocation>
+ <relocation>
+ <pattern>org.xerial.snappy</pattern>
+
<shadedPattern>org.apache.pulsar.shade.org.xerial.snappy</shadedPattern>
+ </relocation>
+ <relocation>
+ <pattern>org.tukaani</pattern>
+
<shadedPattern>org.apache.pulsar.shade.org.tukaani</shadedPattern>
+ </relocation>
</relocations>
<filters>
<filter>
diff --git a/pulsar-client-shaded/pom.xml b/pulsar-client-shaded/pom.xml
index e1e6abf..3df4105 100644
--- a/pulsar-client-shaded/pom.xml
+++ b/pulsar-client-shaded/pom.xml
@@ -86,6 +86,14 @@
<include>org.apache.httpcomponents:httpclient</include>
<include>commons-logging:commons-logging</include>
<include>org.apache.httpcomponents:httpcore</include>
+ <include>org.apache.avro:avro</include>
+ <!-- Avro transitive dependencies-->
+ <include>org.codehaus.jackson:jackson-core-asl</include>
+ <include>org.codehaus.jackson:jackson-mapper-asl</include>
+ <include>com.thoughtworks.paranamer:paranamer</include>
+ <include>org.xerial.snappy:snappy-java</include>
+ <include>org.apache.commons:commons-compress</include>
+ <include>org.tukaani:xz</include>
</includes>
</artifactSet>
<filters>
@@ -167,6 +175,27 @@
<pattern>org.apache.http</pattern>
<shadedPattern>org.apache.pulsar.shade.org.apache.http</shadedPattern>
</relocation>
+ <relocation>
+ <pattern>org.apache.avro</pattern>
+
<shadedPattern>org.apache.pulsar.shade.org.apache.avro</shadedPattern>
+ </relocation>
+ <!-- Avro transitive dependencies-->
+ <relocation>
+ <pattern>org.codehaus.jackson</pattern>
+
<shadedPattern>org.apache.pulsar.shade.org.codehaus.jackson</shadedPattern>
+ </relocation>
+ <relocation>
+ <pattern>com.thoughtworks.paranamer</pattern>
+
<shadedPattern>org.apache.pulsar.shade.com.thoughtworks.paranamer</shadedPattern>
+ </relocation>
+ <relocation>
+ <pattern>org.xerial.snappy</pattern>
+
<shadedPattern>org.apache.pulsar.shade.org.xerial.snappy</shadedPattern>
+ </relocation>
+ <relocation>
+ <pattern>org.tukaani</pattern>
+
<shadedPattern>org.apache.pulsar.shade.org.tukaani</shadedPattern>
+ </relocation>
</relocations>
<transformers>
<transformer
implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer"
/>
diff --git a/pulsar-client/pom.xml b/pulsar-client/pom.xml
index 16ff286..16a7224 100644
--- a/pulsar-client/pom.xml
+++ b/pulsar-client/pom.xml
@@ -92,6 +92,12 @@
</dependency>
<dependency>
+ <groupId>org.apache.avro</groupId>
+ <artifactId>avro</artifactId>
+ <version>${avro.version}</version>
+ </dependency>
+
+ <dependency>
<groupId>com.google.protobuf</groupId>
<artifactId>protobuf-java</artifactId>
<version>${protobuf3.version}</version>
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
new file mode 100644
index 0000000..4bca999
--- /dev/null
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/AvroSchema.java
@@ -0,0 +1,98 @@
+/**
+ * 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 lombok.extern.slf4j.Slf4j;
+import org.apache.avro.io.BinaryEncoder;
+import org.apache.avro.io.DecoderFactory;
+import org.apache.avro.io.EncoderFactory;
+import org.apache.avro.reflect.ReflectData;
+import org.apache.avro.reflect.ReflectDatumReader;
+import org.apache.avro.reflect.ReflectDatumWriter;
+import org.apache.pulsar.client.api.Schema;
+import org.apache.pulsar.client.api.SchemaSerializationException;
+import org.apache.pulsar.common.schema.SchemaInfo;
+import org.apache.pulsar.common.schema.SchemaType;
+
+import java.io.ByteArrayOutputStream;
+import java.io.IOException;
+import java.util.Collections;
+import java.util.Map;
+
+@Slf4j
+public class AvroSchema<T> implements Schema<T> {
+
+ private SchemaInfo schemaInfo;
+ private org.apache.avro.Schema schema;
+ private ReflectDatumWriter<T> datumWriter;
+ private ReflectDatumReader<T> reader;
+ private BinaryEncoder encoder;
+ private ByteArrayOutputStream byteArrayOutputStream;
+
+ public AvroSchema(Class<T> pojo, Map<String, String> properties) {
+ this.schema = ReflectData.AllowNull.get().getSchema(pojo);
+
+ this.schemaInfo = new SchemaInfo();
+ this.schemaInfo.setName("");
+ this.schemaInfo.setProperties(properties);
+ this.schemaInfo.setType(SchemaType.AVRO);
+ this.schemaInfo.setSchema(this.schema.toString().getBytes());
+
+ this.byteArrayOutputStream = new ByteArrayOutputStream();
+ this.encoder =
EncoderFactory.get().binaryEncoder(this.byteArrayOutputStream, this.encoder);
+ this.datumWriter = new ReflectDatumWriter<>(this.schema);
+ this.reader = new ReflectDatumReader<>(this.schema);
+ }
+
+ @Override
+ public byte[] encode(T message) {
+
+ try {
+ datumWriter.write(message, this.encoder);
+ this.encoder.flush();
+ return this.byteArrayOutputStream.toByteArray();
+ } catch (Exception e) {
+ throw new SchemaSerializationException(e);
+ } finally {
+ this.byteArrayOutputStream.reset();
+ }
+ }
+
+ @Override
+ public T decode(byte[] bytes) {
+ try {
+ return reader.read(null, DecoderFactory.get().binaryDecoder(bytes,
null));
+ } catch (IOException e) {
+ throw new SchemaSerializationException(e);
+ }
+ }
+
+ @Override
+ public SchemaInfo getSchemaInfo() {
+ return this.schemaInfo;
+ }
+
+ public static <T> AvroSchema<T> of(Class<T> pojo) {
+ return new AvroSchema<>(pojo, Collections.emptyMap());
+ }
+
+ public static <T> AvroSchema<T> of(Class<T> pojo, Map<String, String>
properties) {
+ return new AvroSchema<>(pojo, properties);
+ }
+}
diff --git
a/pulsar-client/src/test/java/org/apache/pulsar/client/schemas/AvroSchemaTest.java
b/pulsar-client/src/test/java/org/apache/pulsar/client/schemas/AvroSchemaTest.java
new file mode 100644
index 0000000..a8b94de
--- /dev/null
+++
b/pulsar-client/src/test/java/org/apache/pulsar/client/schemas/AvroSchemaTest.java
@@ -0,0 +1,109 @@
+/**
+ * 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.schemas;
+
+import lombok.Data;
+import lombok.EqualsAndHashCode;
+import lombok.ToString;
+import lombok.extern.slf4j.Slf4j;
+import org.apache.avro.Schema;
+import org.apache.pulsar.client.impl.schema.AvroSchema;
+import org.apache.pulsar.common.schema.SchemaType;
+import org.testng.Assert;
+import org.testng.annotations.Test;
+
+@Slf4j
+public class AvroSchemaTest {
+
+ @Data
+ @ToString
+ @EqualsAndHashCode
+ private static class Foo {
+ private String field1;
+ private String field2;
+ private int field3;
+ private Bar field4;
+ }
+
+ @Data
+ @ToString
+ @EqualsAndHashCode
+ private static class Bar {
+ private boolean field1;
+ }
+
+ private static final String SCHEMA_JSON =
"{\"type\":\"record\",\"name\":\"Foo\",\"namespace\":\"org.apache" +
+ ".pulsar.client" +
+
".schemas.AvroSchemaTest$\",\"fields\":[{\"name\":\"field1\",\"type\":[\"null\",\"string\"],"
+
+
"\"default\":null},{\"name\":\"field2\",\"type\":[\"null\",\"string\"],\"default\":null},"
+
+
"{\"name\":\"field3\",\"type\":\"int\"},{\"name\":\"field4\",\"type\":[\"null\",{\"type\":\"record\","
+
+
"\"name\":\"Bar\",\"fields\":[{\"name\":\"field1\",\"type\":\"boolean\"}]}],\"default\":null}]}";
+
+ private static String[] FOO_FIELDS = {
+ "field1",
+ "field2",
+ "field3",
+ "field4"
+ };
+
+ @Test
+ public void testSchema() {
+ AvroSchema<Foo> avroSchema = AvroSchema.of(Foo.class);
+ Assert.assertEquals(avroSchema.getSchemaInfo().getType(),
SchemaType.AVRO);
+ Schema.Parser parser = new Schema.Parser();
+ String schemaJson = new String(avroSchema.getSchemaInfo().getSchema());
+ Assert.assertEquals(schemaJson, SCHEMA_JSON);
+ Schema schema = parser.parse(schemaJson);
+
+ for (String fieldName : FOO_FIELDS) {
+ Schema.Field field = schema.getField(fieldName);
+ Assert.assertNotNull(field);
+
+ if (field.name().equals("field4")) {
+
Assert.assertNotNull(field.schema().getTypes().get(1).getField("field1"));
+ }
+ }
+ }
+
+ @Test
+ public void testEncodeAndDecode() {
+ AvroSchema<Foo> avroSchema = AvroSchema.of(Foo.class, null);
+
+ Foo foo1 = new Foo();
+ foo1.setField1("foo1");
+ foo1.setField2("bar1");
+ foo1.setField4(new Bar());
+
+ Foo foo2 = new Foo();
+ foo2.setField1("foo2");
+ foo2.setField2("bar2");
+
+ byte[] bytes1 = avroSchema.encode(foo1);
+ Assert.assertTrue(bytes1.length > 0);
+
+ byte[] bytes2 = avroSchema.encode(foo2);
+ Assert.assertTrue(bytes2.length > 0);
+
+ Foo object1 = avroSchema.decode(bytes1);
+ Foo object2 = avroSchema.decode(bytes2);
+
+ Assert.assertEquals(object1, foo1);
+ Assert.assertEquals(object2, foo2);
+ }
+}
diff --git
a/pulsar-common/src/main/java/org/apache/pulsar/common/api/Commands.java
b/pulsar-common/src/main/java/org/apache/pulsar/common/api/Commands.java
index 84d0218..1fd38a0 100644
--- a/pulsar-common/src/main/java/org/apache/pulsar/common/api/Commands.java
+++ b/pulsar-common/src/main/java/org/apache/pulsar/common/api/Commands.java
@@ -453,6 +453,10 @@ public class Commands {
return PulsarApi.Schema.Type.String;
case JSON:
return PulsarApi.Schema.Type.Json;
+ case PROTOBUF:
+ return PulsarApi.Schema.Type.Protobuf;
+ case AVRO:
+ return PulsarApi.Schema.Type.Avro;
default:
return PulsarApi.Schema.Type.None;
}
diff --git
a/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java
b/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java
index ce0124b..5a6329b 100644
---
a/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java
+++
b/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java
@@ -319,11 +319,15 @@ public final class PulsarApi {
None(0, 0),
String(1, 1),
Json(2, 2),
+ Protobuf(3, 3),
+ Avro(4, 4),
;
public static final int None_VALUE = 0;
public static final int String_VALUE = 1;
public static final int Json_VALUE = 2;
+ public static final int Protobuf_VALUE = 3;
+ public static final int Avro_VALUE = 4;
public final int getNumber() { return value; }
@@ -333,6 +337,8 @@ public final class PulsarApi {
case 0: return None;
case 1: return String;
case 2: return Json;
+ case 3: return Protobuf;
+ case 4: return Avro;
default: return null;
}
}
diff --git
a/pulsar-common/src/main/java/org/apache/pulsar/common/schema/SchemaType.java
b/pulsar-common/src/main/java/org/apache/pulsar/common/schema/SchemaType.java
index 44d78c9..cbf7c91 100644
---
a/pulsar-common/src/main/java/org/apache/pulsar/common/schema/SchemaType.java
+++
b/pulsar-common/src/main/java/org/apache/pulsar/common/schema/SchemaType.java
@@ -37,5 +37,13 @@ public enum SchemaType {
*/
JSON,
- PROTOBUF
+ /**
+ * Protobuf message encoding and decoding
+ */
+ PROTOBUF,
+
+ /**
+ * Serialize and deserialize via avro
+ */
+ AVRO
}
diff --git a/pulsar-common/src/main/proto/PulsarApi.proto
b/pulsar-common/src/main/proto/PulsarApi.proto
index 3f9de1a..d2ac4fd 100644
--- a/pulsar-common/src/main/proto/PulsarApi.proto
+++ b/pulsar-common/src/main/proto/PulsarApi.proto
@@ -27,6 +27,8 @@ message Schema {
None = 0;
String = 1;
Json = 2;
+ Protobuf = 3;
+ Avro = 4;
}
required string name = 1;
--
To stop receiving notification emails like this one, please contact
[email protected].