Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
56 changes: 56 additions & 0 deletions .github/workflows/external-schemas.yaml
Original file line number Diff line number Diff line change
@@ -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
4 changes: 4 additions & 0 deletions external-schemas/java-demos/avro/.gitignore
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
target/
.idea/
*.iml
.env
22 changes: 22 additions & 0 deletions external-schemas/java-demos/avro/README.md
Original file line number Diff line number Diff line change
@@ -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='<raw-jwt-without-token-prefix>'

# 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.
73 changes: 73 additions & 0 deletions external-schemas/java-demos/avro/pom.xml
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>io.streamnative.examples</groupId>
<artifactId>avro-schema-reference-demo</artifactId>
<version>1.0-SNAPSHOT</version>
<name>Pulsar External Avro Schema Reference Demo</name>
<properties>
<maven.compiler.release>17</maven.compiler.release>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<confluent.version>8.0.0</confluent.version>
<external-schemas.version>1.0.0</external-schemas.version>
<pulsar.version>4.1.0</pulsar.version>
</properties>
<repositories>
<repository>
<id>confluent</id>
<url>https://packages.confluent.io/maven/</url>
</repository>
</repositories>
<dependencies>
<dependency>
<groupId>io.streamnative.schemas.external</groupId>
<artifactId>kafka-schemas</artifactId>
<version>${external-schemas.version}</version>
</dependency>
<dependency>
<groupId>org.apache.pulsar</groupId>
<artifactId>pulsar-client-original</artifactId>
<version>${pulsar.version}</version>
</dependency>
<dependency>
<groupId>io.confluent</groupId>
<artifactId>kafka-avro-serializer</artifactId>
<version>${confluent.version}</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.apache.avro</groupId>
<artifactId>avro-maven-plugin</artifactId>
<version>1.12.0</version>
<executions>
<execution>
<phase>generate-sources</phase>
<goals><goal>schema</goal></goals>
</execution>
</executions>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.8.1</version>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-dependency-plugin</artifactId>
<version>3.7.0</version>
<executions>
<execution>
<id>copy-runtime-dependencies</id>
<phase>package</phase>
<goals><goal>copy-dependencies</goal></goals>
<configuration><includeScope>runtime</includeScope></configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
16 changes: 16 additions & 0 deletions external-schemas/java-demos/avro/src/main/avro/Member.avsc
Original file line number Diff line number Diff line change
@@ -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"}
]
}}
]
}
Original file line number Diff line number Diff line change
@@ -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<String, Object> registryConfigs() {
Map<String, Object> 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;
}
}
Original file line number Diff line number Diff line change
@@ -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<String, Object> 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<Member> 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);
}
}
}
Loading
Loading