From 828b3ef1ed948c75db6ea54904d839d8be7fb31c Mon Sep 17 00:00:00 2001 From: "gaoran_10@126.com" Date: Thu, 1 Oct 2026 02:06:09 +0800 Subject: [PATCH 1/2] [external-schemas] Add Avro external schemas demo --- external-schemas/java-demos/avro/.gitignore | 4 + external-schemas/java-demos/avro/README.md | 22 ++++ external-schemas/java-demos/avro/pom.xml | 73 ++++++++++++ .../java-demos/avro/src/main/avro/Member.avsc | 16 +++ .../examples/avro/DemoConfig.java | 59 ++++++++++ .../examples/avro/SchemaConsume.java | 85 ++++++++++++++ .../examples/avro/SchemaProduce.java | 109 ++++++++++++++++++ 7 files changed, 368 insertions(+) create mode 100644 external-schemas/java-demos/avro/.gitignore create mode 100644 external-schemas/java-demos/avro/README.md create mode 100644 external-schemas/java-demos/avro/pom.xml create mode 100644 external-schemas/java-demos/avro/src/main/avro/Member.avsc create mode 100644 external-schemas/java-demos/avro/src/main/java/io/streamnative/examples/avro/DemoConfig.java create mode 100644 external-schemas/java-demos/avro/src/main/java/io/streamnative/examples/avro/SchemaConsume.java create mode 100644 external-schemas/java-demos/avro/src/main/java/io/streamnative/examples/avro/SchemaProduce.java diff --git a/external-schemas/java-demos/avro/.gitignore b/external-schemas/java-demos/avro/.gitignore new file mode 100644 index 0000000..89ed517 --- /dev/null +++ b/external-schemas/java-demos/avro/.gitignore @@ -0,0 +1,4 @@ +target/ +.idea/ +*.iml +.env diff --git a/external-schemas/java-demos/avro/README.md b/external-schemas/java-demos/avro/README.md new file mode 100644 index 0000000..110ae54 --- /dev/null +++ b/external-schemas/java-demos/avro/README.md @@ -0,0 +1,22 @@ +# Usage + +Requires JDK 17+ and Maven 3.8+. + +```sh +export PULSAR_SERVICE_URL='pulsar+ssl://your-cluster:6651' +export SCHEMA_REGISTRY_URL='https://your-cluster/kafka' +export PULSAR_TOKEN='' + +# Optional +export PULSAR_TOPIC='test-pulsar-external-avro-reference' +export PULSAR_SUBSCRIPTION='pulsar-avro-reference-sub' +export MESSAGE_COUNT=10 + +mvn clean package + +# Send and receive MESSAGE_COUNT messages (default: 10) +java -cp 'target/classes:target/dependency/*' io.streamnative.examples.avro.SchemaProduce +java -cp 'target/classes:target/dependency/*' io.streamnative.examples.avro.SchemaConsume +``` + +In an IDE, import `pom.xml`, set the same environment variables, and run the two main classes. diff --git a/external-schemas/java-demos/avro/pom.xml b/external-schemas/java-demos/avro/pom.xml new file mode 100644 index 0000000..01e3dd4 --- /dev/null +++ b/external-schemas/java-demos/avro/pom.xml @@ -0,0 +1,73 @@ + + + 4.0.0 + io.streamnative.examples + avro-schema-reference-demo + 1.0-SNAPSHOT + Pulsar External Avro Schema Reference Demo + + 17 + UTF-8 + 8.0.0 + 1.0.0 + 4.1.0 + + + + confluent + https://packages.confluent.io/maven/ + + + + + io.streamnative.schemas.external + kafka-schemas + ${external-schemas.version} + + + org.apache.pulsar + pulsar-client-original + ${pulsar.version} + + + io.confluent + kafka-avro-serializer + ${confluent.version} + + + + + + org.apache.avro + avro-maven-plugin + 1.12.0 + + + generate-sources + schema + + + + + org.apache.maven.plugins + maven-compiler-plugin + 3.8.1 + + + org.apache.maven.plugins + maven-dependency-plugin + 3.7.0 + + + copy-runtime-dependencies + package + copy-dependencies + runtime + + + + + + diff --git a/external-schemas/java-demos/avro/src/main/avro/Member.avsc b/external-schemas/java-demos/avro/src/main/avro/Member.avsc new file mode 100644 index 0000000..e20e41b --- /dev/null +++ b/external-schemas/java-demos/avro/src/main/avro/Member.avsc @@ -0,0 +1,16 @@ +{ + "type": "record", + "name": "Member", + "namespace": "io.streamnative.schemas.external.test.avro", + "fields": [ + {"name": "name", "type": "string"}, + {"name": "address", "type": { + "type": "record", + "name": "Address", + "fields": [ + {"name": "city", "type": "string"}, + {"name": "zip", "type": "string"} + ] + }} + ] +} diff --git a/external-schemas/java-demos/avro/src/main/java/io/streamnative/examples/avro/DemoConfig.java b/external-schemas/java-demos/avro/src/main/java/io/streamnative/examples/avro/DemoConfig.java new file mode 100644 index 0000000..3787f08 --- /dev/null +++ b/external-schemas/java-demos/avro/src/main/java/io/streamnative/examples/avro/DemoConfig.java @@ -0,0 +1,59 @@ +package io.streamnative.examples.avro; + +import java.util.HashMap; +import java.util.Map; + +/** Shared environment configuration for the Pulsar producer and consumer. */ +final class DemoConfig { + private DemoConfig() {} + + static String serviceUrl() { + return required("PULSAR_SERVICE_URL"); + } + + static String token() { + return required("PULSAR_TOKEN"); + } + + static Map registryConfigs() { + Map configs = new HashMap<>(); + configs.put("schema.registry.url", required("SCHEMA_REGISTRY_URL")); + configs.put("basic.auth.credentials.source", "USER_INFO"); + configs.put("basic.auth.user.info", "public:token:" + token()); + return configs; + } + + static String topic() { + return optional("PULSAR_TOPIC", "test-pulsar-external-avro-reference"); + } + + static int messageCount() { + String value = optional("MESSAGE_COUNT", "10"); + try { + int count = Integer.parseInt(value); + if (count > 0) { + return count; + } + } catch (NumberFormatException ignored) { + // Report the same validation error for malformed and out-of-range values. + } + throw new IllegalArgumentException("MESSAGE_COUNT must be a positive integer"); + } + + static String subscription() { + return optional("PULSAR_SUBSCRIPTION", "pulsar-avro-reference-sub"); + } + + private static String required(String name) { + String value = System.getenv(name); + if (value == null || value.isBlank()) { + throw new IllegalArgumentException("Set environment variable " + name); + } + return value; + } + + private static String optional(String name, String fallback) { + String value = System.getenv(name); + return value == null || value.isBlank() ? fallback : value; + } +} diff --git a/external-schemas/java-demos/avro/src/main/java/io/streamnative/examples/avro/SchemaConsume.java b/external-schemas/java-demos/avro/src/main/java/io/streamnative/examples/avro/SchemaConsume.java new file mode 100644 index 0000000..5e22499 --- /dev/null +++ b/external-schemas/java-demos/avro/src/main/java/io/streamnative/examples/avro/SchemaConsume.java @@ -0,0 +1,85 @@ +/** + * Copyright StreamNative, Inc. and contributors + * + * Licensed 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 io.streamnative.examples.avro; + +import java.nio.ByteBuffer; +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.TimeUnit; + +import io.confluent.kafka.schemaregistry.avro.AvroSchema; +import io.confluent.kafka.serializers.schema.id.SchemaId; +import io.streamnative.schemas.external.KafkaSchemaFactory; +import io.streamnative.schemas.external.test.avro.Member; +import org.apache.pulsar.client.api.AuthenticationFactory; +import org.apache.pulsar.client.api.Consumer; +import org.apache.pulsar.client.api.PulsarClient; +import org.apache.pulsar.client.api.SubscriptionInitialPosition; + +public class SchemaConsume { + + public static void main(String[] args) throws Exception { + int count = DemoConfig.messageCount(); + consume(DemoConfig.serviceUrl(), DemoConfig.registryConfigs(), DemoConfig.token(), + DemoConfig.topic(), DemoConfig.subscription(), count); + } + + static void consume( + String serviceUrl, + Map registryConfigs, + String jwtToken, + String topic, + String subscription, + int count) + throws Exception { + var consumerConfigs = new HashMap<>(registryConfigs); + consumerConfigs.put("auto.register.schemas", false); + // The EXTERNAL schema adapter resolves the Registry ID carried in Pulsar message metadata. + var externalSchema = new KafkaSchemaFactory(consumerConfigs).avro(Member.class); + var clientBuilder = PulsarClient.builder().serviceUrl(serviceUrl); + if (jwtToken != null && !jwtToken.isBlank()) { + clientBuilder.authentication(AuthenticationFactory.token(jwtToken)); + } + try (PulsarClient client = clientBuilder.build(); + Consumer consumer = + client.newConsumer(externalSchema) + .topic(topic) + .subscriptionName(subscription) + .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest) + .subscribe()) { + for (int received = 0; received < count; received++) { + var message = consumer.receive(60, TimeUnit.SECONDS); + if (message == null) { + throw new IllegalStateException(String.format( + "Timed out after 60 seconds waiting for the next message on %s; received %d/%d", + topic, received, count)); + } + var schemaId = new SchemaId(AvroSchema.TYPE); + schemaId.fromBytes(ByteBuffer.wrap(message.getSchemaId().orElseThrow())); + Member member = message.getValue(); + if (member == null || member.getAddress() == null) { + throw new IllegalStateException("Expected Member with a referenced Address record"); + } + System.out.printf( + "Received %d/%d: Member{name=%s, address={city=%s, zip=%s}} using external schema ID %d from %s, messageId=%s%n", + received + 1, count, member.getName(), member.getAddress().getCity(), + member.getAddress().getZip(), schemaId.getId(), message.getTopicName(), message.getMessageId()); + consumer.acknowledge(message); + } + System.out.printf("Finished receiving and acknowledging %d messages.%n", count); + } + } +} diff --git a/external-schemas/java-demos/avro/src/main/java/io/streamnative/examples/avro/SchemaProduce.java b/external-schemas/java-demos/avro/src/main/java/io/streamnative/examples/avro/SchemaProduce.java new file mode 100644 index 0000000..7962013 --- /dev/null +++ b/external-schemas/java-demos/avro/src/main/java/io/streamnative/examples/avro/SchemaProduce.java @@ -0,0 +1,109 @@ +/** + * Copyright StreamNative, Inc. and contributors + * + * Licensed 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 io.streamnative.examples.avro; + +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import io.confluent.kafka.schemaregistry.avro.AvroSchema; +import io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient; +import io.confluent.kafka.schemaregistry.client.rest.entities.SchemaReference; +import io.streamnative.schemas.external.KafkaSchemaFactory; +import io.streamnative.schemas.external.test.avro.Address; +import io.streamnative.schemas.external.test.avro.Member; +import org.apache.pulsar.client.api.AuthenticationFactory; +import org.apache.pulsar.client.api.Producer; +import org.apache.pulsar.client.api.PulsarClient; + +public class SchemaProduce { + + public static void main(String[] args) throws Exception { + int count = DemoConfig.messageCount(); + produce(DemoConfig.serviceUrl(), DemoConfig.registryConfigs(), DemoConfig.token(), + DemoConfig.topic(), count); + } + + static int produce( + String serviceUrl, Map registryConfigs, String jwtToken, String topic, int count) + throws Exception { + String registryUrl = (String) registryConfigs.get("schema.registry.url"); + String addressName = Address.getClassSchema().getFullName(); + String addressSubject = topic + "-address"; + String memberSubject = topic + "-value"; + AvroSchema addressSchema = new AvroSchema(Address.getClassSchema()); + String memberDefinition = + """ + { + "type": "record", + "name": "Member", + "namespace": "io.streamnative.schemas.external.test.avro", + "fields": [ + {"name": "name", "type": "string"}, + {"name": "address", "type": "io.streamnative.schemas.external.test.avro.Address"} + ] + } + """; + try (var registry = new CachedSchemaRegistryClient(registryUrl, 10, registryConfigs)) { + int addressId = registry.register(addressSubject, addressSchema); + int version = registry.getVersion(addressSubject, addressSchema); + var reference = new SchemaReference(addressName, addressSubject, version); + var memberSchema = + new AvroSchema( + memberDefinition, + List.of(reference), + Map.of(addressName, addressSchema.canonicalString()), + null); + int memberId = registry.register(memberSubject, memberSchema); + System.out.printf( + "Registered Address: subject=%s, id=%d, version=%d%n", + addressSubject, addressId, version); + System.out.printf( + "Registered Member: subject=%s, id=%d, references=%s%n", + memberSubject, memberId, memberSchema.references()); + + var producerConfigs = new HashMap<>(registryConfigs); + producerConfigs.put("auto.register.schemas", false); + producerConfigs.put("use.schema.id", memberId); + producerConfigs.put("id.compatibility.strict", false); + var externalSchema = new KafkaSchemaFactory(producerConfigs).avro(Member.class); + var clientBuilder = PulsarClient.builder().serviceUrl(serviceUrl); + if (jwtToken != null && !jwtToken.isBlank()) { + // Pulsar token authentication accepts the raw JWT, without the Kafka prefix. + clientBuilder.authentication(AuthenticationFactory.token(jwtToken)); + } + var memberVersions = List.copyOf(registry.getAllVersions(memberSubject)); + var addressVersions = List.copyOf(registry.getAllVersions(addressSubject)); + try (PulsarClient client = clientBuilder.build(); + Producer producer = + client.newProducer(externalSchema).topic(topic).create()) { + for (int i = 0; i < count; i++) { + Member member = new Member("jwt-sr-" + (i + 1), new Address("Shanghai", "200000")); + var messageId = producer.send(member); + System.out.printf( + "Sent %d/%d: %s using external schema ID %d to %s, messageId=%s%n", + i + 1, count, member, memberId, producer.getTopic(), messageId); + } + System.out.printf("Finished sending %d messages.%n", count); + } + if (!registry.getAllVersions(memberSubject).equals(memberVersions) + || !registry.getAllVersions(addressSubject).equals(addressVersions)) { + throw new IllegalStateException("Producing unexpectedly created a schema version"); + } + return memberId; + } + } +} From 5cebffde91af66615c72b8f2bf709c81001a53bc Mon Sep 17 00:00:00 2001 From: "gaoran_10@126.com" Date: Thu, 1 Oct 2026 02:11:10 +0800 Subject: [PATCH 2/2] add --- .github/workflows/external-schemas.yaml | 56 +++++++++++++++++++++++++ 1 file changed, 56 insertions(+) create mode 100644 .github/workflows/external-schemas.yaml diff --git a/.github/workflows/external-schemas.yaml b/.github/workflows/external-schemas.yaml new file mode 100644 index 0000000..ffc9c67 --- /dev/null +++ b/.github/workflows/external-schemas.yaml @@ -0,0 +1,56 @@ +# +# Licensed 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. +# + +name: External Schemas + +on: + pull_request: + branches: + - master + paths: + - 'external-schemas/**' + - '.github/workflows/external-schemas.yaml' + push: + branches: + - master + paths: + - 'external-schemas/**' + - '.github/workflows/external-schemas.yaml' + workflow_dispatch: + +permissions: + contents: read + +jobs: + avro-compile: + name: Compile Avro demo + runs-on: ubuntu-latest + timeout-minutes: 15 + defaults: + run: + working-directory: external-schemas/java-demos/avro + steps: + - name: Checkout + uses: actions/checkout@v4 + + - name: Set up JDK 17 + uses: actions/setup-java@v4 + with: + distribution: temurin + java-version: '17' + cache: maven + cache-dependency-path: external-schemas/java-demos/avro/pom.xml + + - name: Compile external schema demo + run: mvn --batch-mode --no-transfer-progress clean compile