mirror of
https://github.com/wahyd4/schema-registry-transfer-smt.git
synced 2026-08-09 04:25:59 +10:00
- handle case when schema info is present but value is null
- not skipping record in that case but omitting schema re-registration
This commit is contained in:
committed by
Jordan Moore
parent
443ae8d77c
commit
a40d8a97dc
@@ -124,16 +124,18 @@ public class SchemaRegistryTransfer<R extends ConnectRecord<R>> implements Trans
|
||||
|
||||
Object updatedValue;
|
||||
Optional<Integer> destValueSchemaId;
|
||||
if (value == null) {
|
||||
throw new ConnectException("Unable to extract schema information from null record value.");
|
||||
}
|
||||
if ((valueSchema != null && valueSchema.type() == Schema.BYTES_SCHEMA.type()) ||
|
||||
value instanceof byte[]) {
|
||||
ByteBuffer b = ByteBuffer.wrap((byte[]) value);
|
||||
destValueSchemaId = copySchema(b, topic, false);
|
||||
b.putInt(1, destValueSchemaId.orElseThrow(()
|
||||
-> new ConnectException("Transform failed. Unable to update record schema id. (isKey=false)")));
|
||||
updatedValue = b.array();
|
||||
if(value!=null) {
|
||||
ByteBuffer b = ByteBuffer.wrap((byte[]) value);
|
||||
destValueSchemaId = copySchema(b, topic, false);
|
||||
b.putInt(1, destValueSchemaId.orElseThrow(()
|
||||
-> new ConnectException("Transform failed. Unable to update record schema id. (isKey=false)")));
|
||||
updatedValue = b.array();
|
||||
} else {
|
||||
log.trace("Cannot extract schema details from null-value record.");
|
||||
updatedValue = value;
|
||||
}
|
||||
} else {
|
||||
throw new ConnectException("Transform failed. Record value does not have a byte[] schema.");
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user