mirror of
https://github.com/wahyd4/schema-registry-transfer-smt.git
synced 2026-08-09 04:25:59 +10:00
Add tests that copy multiple schemas
This commit is contained in:
@@ -5,6 +5,8 @@ import static java.net.HttpURLConnection.HTTP_NOT_FOUND;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.List;
|
||||
import java.util.stream.Collectors;
|
||||
import java.util.stream.StreamSupport;
|
||||
|
||||
import org.apache.avro.Schema;
|
||||
import org.junit.jupiter.api.extension.AfterEachCallback;
|
||||
@@ -17,6 +19,7 @@ import io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient;
|
||||
import io.confluent.kafka.schemaregistry.client.MockSchemaRegistryClient;
|
||||
import io.confluent.kafka.schemaregistry.client.SchemaMetadata;
|
||||
import io.confluent.kafka.schemaregistry.client.SchemaRegistryClient;
|
||||
import io.confluent.kafka.schemaregistry.client.rest.entities.Config;
|
||||
import io.confluent.kafka.schemaregistry.client.rest.entities.SchemaString;
|
||||
import io.confluent.kafka.schemaregistry.client.rest.entities.requests.RegisterSchemaRequest;
|
||||
import io.confluent.kafka.schemaregistry.client.rest.entities.requests.RegisterSchemaResponse;
|
||||
@@ -62,19 +65,22 @@ import com.google.common.collect.Iterables;
|
||||
* assertThat(retrievedSchema).isEqualTo(keySchema);
|
||||
* }
|
||||
* }</code></pre>
|
||||
*
|
||||
* <p>
|
||||
* To retrieve the url of the schema registry for a Kafka Streams config, please use {@link #getUrl()}
|
||||
*/
|
||||
public class SchemaRegistryMock implements BeforeEachCallback, AfterEachCallback {
|
||||
private static final String SCHEMA_REGISTRATION_PATTERN = "/subjects/[^/]+/versions";
|
||||
private static final String SCHEMA_BY_ID_PATTERN = "/schemas/ids/";
|
||||
private static final String CONFIG_PATTERN = "/config";
|
||||
private static final int IDENTITY_MAP_CAPACITY = 1000;
|
||||
private final ListVersionsHandler listVersionsHandler = new ListVersionsHandler();
|
||||
private final GetVersionHandler getVersionHandler = new GetVersionHandler();
|
||||
private final AutoRegistrationHandler autoRegistrationHandler = new AutoRegistrationHandler();
|
||||
private final GetConfigHandler getConfigHandler = new GetConfigHandler();
|
||||
private final WireMockServer mockSchemaRegistry = new WireMockServer(
|
||||
WireMockConfiguration.wireMockConfig().dynamicPort()
|
||||
.extensions(this.autoRegistrationHandler, this.listVersionsHandler, this.getVersionHandler));
|
||||
WireMockConfiguration.wireMockConfig().dynamicPort().extensions(
|
||||
this.autoRegistrationHandler, this.listVersionsHandler, this.getVersionHandler,
|
||||
this.getConfigHandler));
|
||||
private final SchemaRegistryClient schemaRegistryClient = new MockSchemaRegistryClient();
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(SchemaRegistryMock.class);
|
||||
@@ -93,6 +99,8 @@ public class SchemaRegistryMock implements BeforeEachCallback, AfterEachCallback
|
||||
.willReturn(WireMock.aResponse().withTransformers(this.autoRegistrationHandler.getName())));
|
||||
this.mockSchemaRegistry.stubFor(WireMock.get(WireMock.urlPathMatching(SCHEMA_REGISTRATION_PATTERN + "/(?:latest|\\d+)"))
|
||||
.willReturn(WireMock.aResponse().withTransformers(this.getVersionHandler.getName())));
|
||||
this.mockSchemaRegistry.stubFor(WireMock.get(WireMock.urlPathMatching(CONFIG_PATTERN))
|
||||
.willReturn(WireMock.aResponse().withTransformers(this.getConfigHandler.getName())));
|
||||
this.mockSchemaRegistry.stubFor(WireMock.get(WireMock.urlPathMatching(SCHEMA_BY_ID_PATTERN + "\\d+"))
|
||||
.willReturn(WireMock.aResponse().withStatus(HTTP_NOT_FOUND)));
|
||||
}
|
||||
@@ -131,7 +139,7 @@ public class SchemaRegistryMock implements BeforeEachCallback, AfterEachCallback
|
||||
try {
|
||||
if (version instanceof String && version.equals("latest")) {
|
||||
return this.schemaRegistryClient.getLatestSchemaMetadata(subject);
|
||||
} else if (version instanceof Number){
|
||||
} else if (version instanceof Number) {
|
||||
return this.schemaRegistryClient.getSchemaMetadata(subject, ((Number) version).intValue());
|
||||
} else {
|
||||
throw new IllegalArgumentException("Only 'latest' or integer versions are allowed");
|
||||
@@ -141,6 +149,19 @@ public class SchemaRegistryMock implements BeforeEachCallback, AfterEachCallback
|
||||
}
|
||||
}
|
||||
|
||||
private String getCompatibility(String subject) {
|
||||
if (subject == null) {
|
||||
log.debug("Requesting registry base compatibility");
|
||||
} else {
|
||||
log.debug("Requesting compatibility for subject {}", subject);
|
||||
}
|
||||
try {
|
||||
return this.schemaRegistryClient.getCompatibility(subject);
|
||||
} catch (IOException | RestClientException e) {
|
||||
throw new IllegalStateException("Internal error in mock schema registry client", e);
|
||||
}
|
||||
}
|
||||
|
||||
public SchemaRegistryClient getSchemaRegistryClient() {
|
||||
return new CachedSchemaRegistryClient(this.getUrl(), IDENTITY_MAP_CAPACITY);
|
||||
}
|
||||
@@ -224,4 +245,29 @@ public class SchemaRegistryMock implements BeforeEachCallback, AfterEachCallback
|
||||
}
|
||||
}
|
||||
|
||||
private class GetConfigHandler extends SubjectsVersioHandler {
|
||||
|
||||
@Override
|
||||
protected String getSubject(Request request) {
|
||||
List<String> parts =
|
||||
StreamSupport.stream(this.urlSplitter.split(request.getUrl()).spliterator(), false)
|
||||
.collect(Collectors.toList());
|
||||
|
||||
// return null when this is just /config
|
||||
return parts.size() < 2 ? null : parts.get(1);
|
||||
}
|
||||
|
||||
@Override
|
||||
public ResponseDefinition transform(final Request request, final ResponseDefinition responseDefinition,
|
||||
final FileSource files, final Parameters parameters) {
|
||||
Config config = new Config(SchemaRegistryMock.this.getCompatibility(getSubject(request)));
|
||||
return ResponseDefinitionBuilder.jsonResponse(config);
|
||||
}
|
||||
|
||||
@Override
|
||||
public String getName() {
|
||||
return GetConfigHandler.class.getSimpleName();
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -14,7 +14,9 @@ import java.util.Map;
|
||||
import java.util.UUID;
|
||||
|
||||
import org.apache.avro.SchemaBuilder;
|
||||
import org.apache.avro.generic.GenericData;
|
||||
import org.apache.avro.generic.GenericDatumWriter;
|
||||
import org.apache.avro.generic.GenericRecordBuilder;
|
||||
import org.apache.avro.io.BinaryEncoder;
|
||||
import org.apache.avro.io.DatumWriter;
|
||||
import org.apache.avro.io.EncoderFactory;
|
||||
@@ -24,6 +26,7 @@ import org.apache.kafka.connect.data.Schema;
|
||||
import org.apache.kafka.connect.errors.ConnectException;
|
||||
import org.apache.kafka.connect.source.SourceRecord;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Disabled;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.extension.RegisterExtension;
|
||||
import org.slf4j.Logger;
|
||||
@@ -32,6 +35,11 @@ import org.slf4j.LoggerFactory;
|
||||
import io.confluent.kafka.schemaregistry.client.SchemaMetadata;
|
||||
import io.confluent.kafka.schemaregistry.client.SchemaRegistryClient;
|
||||
import io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException;
|
||||
import io.confluent.kafka.serializers.NonRecordContainer;
|
||||
|
||||
import static org.apache.avro.Schema.Type.INT;
|
||||
import static org.apache.avro.Schema.Type.BOOLEAN;
|
||||
import static org.apache.avro.Schema.Type.STRING;
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
public class TransformTest {
|
||||
@@ -41,7 +49,21 @@ public class TransformTest {
|
||||
public static final String TOPIC = TransformTest.class.getSimpleName();
|
||||
|
||||
private static final byte MAGIC_BYTE = (byte) 0x0;
|
||||
private static final int AVRO_CONTENT_OFFSET = 5;
|
||||
public static final int ID_SIZE = Integer.SIZE / Byte.SIZE;
|
||||
private static final int AVRO_CONTENT_OFFSET = 1 + ID_SIZE;
|
||||
public static final org.apache.avro.Schema INT_SCHEMA = org.apache.avro.Schema.create(INT);
|
||||
public static final org.apache.avro.Schema STRING_SCHEMA = org.apache.avro.Schema.create(STRING);
|
||||
public static final org.apache.avro.Schema BOOLEAN_SCHEMA = org.apache.avro.Schema.create(BOOLEAN);
|
||||
public static final org.apache.avro.Schema NAME_SCHEMA = SchemaBuilder.record("FullName")
|
||||
.namespace("cricket.jmoore.kafka.connect.transforms").fields()
|
||||
.requiredString("first")
|
||||
.requiredString("last")
|
||||
.endRecord();
|
||||
public static final org.apache.avro.Schema NAME_SCHEMA_ALIASED = SchemaBuilder.record("FullName")
|
||||
.namespace("cricket.jmoore.kafka.connect.transforms").fields()
|
||||
.requiredString("first")
|
||||
.name("surname").aliases("last").type().stringType().noDefault()
|
||||
.endRecord();
|
||||
|
||||
@RegisterExtension
|
||||
final SchemaRegistryMock sourceSchemaRegistry = new SchemaRegistryMock();
|
||||
@@ -57,6 +79,10 @@ public class TransformTest {
|
||||
return new SourceRecord(null, null, TOPIC, keySchema, key, valueSchema, value);
|
||||
}
|
||||
|
||||
private ConnectRecord createRecord(byte[] key, byte[] value) {
|
||||
return createRecord(Schema.OPTIONAL_BYTES_SCHEMA, key, Schema.OPTIONAL_BYTES_SCHEMA, value);
|
||||
}
|
||||
|
||||
private Map<String, Object> getRequiredTransformConfigs() {
|
||||
Map<String, Object> configs = new HashMap<>();
|
||||
configs.put(ConfigName.SRC_SCHEMA_REGISTRY_URL, sourceSchemaRegistry.getUrl());
|
||||
@@ -79,12 +105,16 @@ public class TransformTest {
|
||||
ByteArrayOutputStream out = new ByteArrayOutputStream();
|
||||
|
||||
out.write(MAGIC_BYTE);
|
||||
out.write(ByteBuffer.allocate(Integer.SIZE / Byte.SIZE).putInt(sourceId).array());
|
||||
out.write(ByteBuffer.allocate(ID_SIZE).putInt(sourceId).array());
|
||||
|
||||
EncoderFactory encoderFactory = EncoderFactory.get();
|
||||
BinaryEncoder encoder = encoderFactory.directBinaryEncoder(out, null);
|
||||
Object
|
||||
value =
|
||||
datum instanceof NonRecordContainer ? ((NonRecordContainer) datum).getValue()
|
||||
: datum;
|
||||
DatumWriter<Object> writer = new GenericDatumWriter<>(schema);
|
||||
writer.write(datum, encoder);
|
||||
writer.write(value, encoder);
|
||||
encoder.flush();
|
||||
|
||||
return out;
|
||||
@@ -168,10 +198,10 @@ public class TransformTest {
|
||||
|
||||
// Create bogus schema in destination so that source and destination ids differ
|
||||
log.info("Registering schema in destination registry");
|
||||
destSchemaRegistry.registerSchema(UUID.randomUUID().toString(), true, SchemaBuilder.builder().intType());
|
||||
destSchemaRegistry.registerSchema(UUID.randomUUID().toString(), true, INT_SCHEMA);
|
||||
|
||||
// Create new schema for source registry
|
||||
org.apache.avro.Schema schema = SchemaBuilder.builder().stringType();
|
||||
org.apache.avro.Schema schema = STRING_SCHEMA;
|
||||
log.info("Registering schema in source registry");
|
||||
int sourceId = sourceSchemaRegistry.registerSchema(TOPIC, true, schema);
|
||||
final String subject = TOPIC + "-key";
|
||||
@@ -189,7 +219,7 @@ public class TransformTest {
|
||||
try {
|
||||
ByteArrayOutputStream out = encodeAvroObject(schema, sourceId, "hello, world");
|
||||
|
||||
ConnectRecord record = createRecord(Schema.STRING_SCHEMA, out.toByteArray(), null, null);
|
||||
ConnectRecord record = createRecord(Schema.OPTIONAL_BYTES_SCHEMA, out.toByteArray(), null, null);
|
||||
|
||||
// check the destination has no versions for this subject
|
||||
SchemaRegistryClient destClient = destSchemaRegistry.getSchemaRegistryClient();
|
||||
@@ -231,10 +261,10 @@ public class TransformTest {
|
||||
|
||||
// Create bogus schema in destination so that source and destination ids differ
|
||||
log.info("Registering schema in destination registry");
|
||||
destSchemaRegistry.registerSchema(UUID.randomUUID().toString(), false, SchemaBuilder.builder().intType());
|
||||
destSchemaRegistry.registerSchema(UUID.randomUUID().toString(), false, INT_SCHEMA);
|
||||
|
||||
// Create new schema for source registry
|
||||
org.apache.avro.Schema schema = SchemaBuilder.builder().stringType();
|
||||
org.apache.avro.Schema schema = STRING_SCHEMA;
|
||||
log.info("Registering schema in source registry");
|
||||
int sourceId = sourceSchemaRegistry.registerSchema(TOPIC, false, schema);
|
||||
final String subject = TOPIC + "-value";
|
||||
@@ -255,9 +285,8 @@ public class TransformTest {
|
||||
try {
|
||||
ByteArrayOutputStream out = encodeAvroObject(schema, sourceId, "hello, world");
|
||||
|
||||
// Giving the key an optional bytes schema so transform doesn't error-out
|
||||
value = out.toByteArray();
|
||||
ConnectRecord record = createRecord(Schema.OPTIONAL_BYTES_SCHEMA, null, Schema.STRING_SCHEMA, value);
|
||||
ConnectRecord record = createRecord(null, value);
|
||||
|
||||
// check the destination has no versions for this subject
|
||||
SchemaRegistryClient destClient = destSchemaRegistry.getSchemaRegistryClient();
|
||||
@@ -313,11 +342,11 @@ public class TransformTest {
|
||||
|
||||
// Create bogus schema in destination so that source and destination ids differ
|
||||
log.info("Registering schema in destination registry");
|
||||
destSchemaRegistry.registerSchema(UUID.randomUUID().toString(), false, SchemaBuilder.builder().booleanType());
|
||||
destSchemaRegistry.registerSchema(UUID.randomUUID().toString(), false, BOOLEAN_SCHEMA);
|
||||
|
||||
// Create new schemas for source registry
|
||||
org.apache.avro.Schema keySchema = SchemaBuilder.builder().intType();
|
||||
org.apache.avro.Schema valueSchema = SchemaBuilder.builder().stringType();
|
||||
org.apache.avro.Schema keySchema = INT_SCHEMA;
|
||||
org.apache.avro.Schema valueSchema = STRING_SCHEMA;
|
||||
log.info("Registering schemas in source registry");
|
||||
int sourceKeyId = sourceSchemaRegistry.registerSchema(TOPIC, true, keySchema);
|
||||
final String keySubject = TOPIC + "-key";
|
||||
@@ -349,7 +378,7 @@ public class TransformTest {
|
||||
|
||||
key = keyStream.toByteArray();
|
||||
value = valueStream.toByteArray();
|
||||
ConnectRecord record = createRecord(Schema.INT32_SCHEMA, key, Schema.STRING_SCHEMA, value);
|
||||
ConnectRecord record = createRecord(key, value);
|
||||
|
||||
// check the destination has no versions for this subject
|
||||
SchemaRegistryClient destClient = destSchemaRegistry.getSchemaRegistryClient();
|
||||
@@ -438,4 +467,139 @@ public class TransformTest {
|
||||
assertEquals(record.valueSchema(), appliedRecord.valueSchema(), "value schema unchanged");
|
||||
assertNull(appliedRecord.value());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEvolvingValueSchemaTransfer() {
|
||||
configure(true);
|
||||
|
||||
// Create bogus schema in destination so that source and destination ids differ
|
||||
log.info("Registering schema in destination registry");
|
||||
destSchemaRegistry.registerSchema(UUID.randomUUID().toString(), false, INT_SCHEMA);
|
||||
|
||||
log.info("Registering schema in source registry");
|
||||
int sourceId = sourceSchemaRegistry.registerSchema(TOPIC, false, NAME_SCHEMA);
|
||||
int nextSourceId = sourceSchemaRegistry.registerSchema(TOPIC, false, NAME_SCHEMA_ALIASED);
|
||||
final String subject = TOPIC + "-value";
|
||||
assertEquals(1, sourceId, "An empty registry starts at id=1");
|
||||
assertEquals(2, nextSourceId, "The next schema is id=2");
|
||||
|
||||
SchemaRegistryClient sourceClient = sourceSchemaRegistry.getSchemaRegistryClient();
|
||||
int numSourceVersions = 0;
|
||||
try {
|
||||
numSourceVersions = sourceClient.getAllVersions(subject).size();
|
||||
assertEquals(2, numSourceVersions, "the source registry subject contains the pre-registered schema");
|
||||
} catch (IOException | RestClientException e) {
|
||||
fail(e);
|
||||
}
|
||||
|
||||
try {
|
||||
GenericData.Record record1 = new GenericRecordBuilder(NAME_SCHEMA)
|
||||
.set("first", "fname")
|
||||
.set("last", "lname")
|
||||
.build();
|
||||
ByteArrayOutputStream out = encodeAvroObject(NAME_SCHEMA, sourceId, record1);
|
||||
|
||||
byte[] value = out.toByteArray();
|
||||
ConnectRecord record = createRecord(null, value);
|
||||
|
||||
GenericData.Record record2 = new GenericRecordBuilder(NAME_SCHEMA_ALIASED)
|
||||
.set("first", "fname")
|
||||
.set("surname", "lname")
|
||||
.build();
|
||||
out = encodeAvroObject(NAME_SCHEMA_ALIASED, nextSourceId, record2);
|
||||
|
||||
byte[] nextValue = out.toByteArray();
|
||||
ConnectRecord nextRecord = createRecord(null, nextValue);
|
||||
|
||||
// check the destination has no versions for this subject
|
||||
SchemaRegistryClient destClient = destSchemaRegistry.getSchemaRegistryClient();
|
||||
List<Integer> destVersions = destClient.getAllVersions(subject);
|
||||
assertTrue(destVersions.isEmpty(), "the destination registry starts empty");
|
||||
|
||||
// The transform will pass for key and value with byte schemas
|
||||
log.info("applying transformation");
|
||||
assertDoesNotThrow(() -> smt.apply(record));
|
||||
|
||||
// check the value schema was copied, and the destination now has some version
|
||||
destVersions = destClient.getAllVersions(subject);
|
||||
assertEquals(1, destVersions.size(),
|
||||
"the destination registry has been updated with first schema");
|
||||
|
||||
log.info("applying transformation");
|
||||
assertDoesNotThrow(() -> smt.apply(nextRecord));
|
||||
|
||||
destVersions = destClient.getAllVersions(subject);
|
||||
assertEquals(numSourceVersions, destVersions.size(),
|
||||
"the destination registry has been updated with the second schema");
|
||||
|
||||
} catch (IOException | RestClientException e) {
|
||||
fail(e);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@Disabled("TODO: Find scenario where a backwards compatible change cannot be undone")
|
||||
public void testIncompatibleEvolvingValueSchemaTransfer() {
|
||||
configure(true);
|
||||
|
||||
// Create bogus schema in destination so that source and destination ids differ
|
||||
log.info("Registering schema in destination registry");
|
||||
destSchemaRegistry.registerSchema(UUID.randomUUID().toString(), false, INT_SCHEMA);
|
||||
|
||||
// Create new schema for source registry
|
||||
log.info("Registering schema in source registry");
|
||||
|
||||
// TODO: Figure out what these should be, where if order is flipped, destination will not accept
|
||||
org.apache.avro.Schema schema = null;
|
||||
org.apache.avro.Schema nextSchema = null;
|
||||
|
||||
int sourceId = sourceSchemaRegistry.registerSchema(TOPIC, false, schema);
|
||||
int nextSourceId = sourceSchemaRegistry.registerSchema(TOPIC, false, nextSchema);
|
||||
final String subject = TOPIC + "-value";
|
||||
assertEquals(1, sourceId, "An empty registry starts at id=1");
|
||||
assertEquals(2, nextSourceId, "The next schema is id=2");
|
||||
|
||||
SchemaRegistryClient sourceClient = sourceSchemaRegistry.getSchemaRegistryClient();
|
||||
int numSourceVersions = 0;
|
||||
try {
|
||||
numSourceVersions = sourceClient.getAllVersions(subject).size();
|
||||
assertEquals(2, numSourceVersions, "the source registry subject contains the pre-registered schema");
|
||||
} catch (IOException | RestClientException e) {
|
||||
fail(e);
|
||||
}
|
||||
|
||||
try {
|
||||
// TODO: Depending on schemas above, then build Avro records for them
|
||||
// ensure second id is encoded first
|
||||
ByteArrayOutputStream out = encodeAvroObject(nextSchema, nextSourceId, null);
|
||||
|
||||
byte[] value = out.toByteArray();
|
||||
ConnectRecord record = createRecord(null, value);
|
||||
|
||||
out = encodeAvroObject(schema, sourceId, null);
|
||||
|
||||
byte[] nextValue = out.toByteArray();
|
||||
ConnectRecord nextRecord = createRecord(null, nextValue);
|
||||
|
||||
// check the destination has no versions for this subject
|
||||
SchemaRegistryClient destClient = destSchemaRegistry.getSchemaRegistryClient();
|
||||
List<Integer> destVersions = destClient.getAllVersions(subject);
|
||||
assertTrue(destVersions.isEmpty(), "the destination registry starts empty");
|
||||
|
||||
// The transform will pass for key and value with byte schemas
|
||||
log.info("applying transformation");
|
||||
assertDoesNotThrow(() -> smt.apply(record));
|
||||
|
||||
// check the value schema was copied, and the destination now has some version
|
||||
destVersions = destClient.getAllVersions(subject);
|
||||
assertEquals(1, destVersions.size(),
|
||||
"the destination registry has been updated with first schema");
|
||||
|
||||
log.info("applying transformation");
|
||||
assertThrows(ConnectException.class, () -> smt.apply(nextRecord));
|
||||
|
||||
} catch (IOException | RestClientException e) {
|
||||
fail(e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user