mirror of
https://github.com/nosqlbench/nosqlbench.git
synced 2025-02-25 18:55:28 -06:00
Add caching to Avro schema instantiation
This commit is contained in:
parent
4948a4101f
commit
a5b0b3cb2d
@ -37,6 +37,7 @@ import java.util.Arrays;
|
|||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.Optional;
|
import java.util.Optional;
|
||||||
import java.util.Set;
|
import java.util.Set;
|
||||||
|
import java.util.concurrent.ConcurrentHashMap;
|
||||||
import java.util.stream.Collectors;
|
import java.util.stream.Collectors;
|
||||||
import java.util.stream.Stream;
|
import java.util.stream.Stream;
|
||||||
|
|
||||||
@ -364,30 +365,31 @@ public class PulsarActivityUtil {
|
|||||||
public static boolean isAutoConsumeSchemaTypeStr(String typeStr) {
|
public static boolean isAutoConsumeSchemaTypeStr(String typeStr) {
|
||||||
return typeStr.equalsIgnoreCase("AUTO_CONSUME");
|
return typeStr.equalsIgnoreCase("AUTO_CONSUME");
|
||||||
}
|
}
|
||||||
public static Schema<?> getAvroSchema(String typeStr, String definitionStr) {
|
|
||||||
String schemaDefinitionStr = definitionStr;
|
|
||||||
String filePrefix = "file://";
|
|
||||||
Schema<?> schema;
|
|
||||||
|
|
||||||
|
|
||||||
|
private static final Map<String, Schema<?>> AVRO_SCHEMA_CACHE = new ConcurrentHashMap<>();
|
||||||
|
|
||||||
|
public static Schema<?> getAvroSchema(String typeStr, final String definitionStr) {
|
||||||
// Check if payloadStr points to a file (e.g. "file:///path/to/a/file")
|
// Check if payloadStr points to a file (e.g. "file:///path/to/a/file")
|
||||||
if (isAvroSchemaTypeStr(typeStr)) {
|
if (isAvroSchemaTypeStr(typeStr)) {
|
||||||
if (StringUtils.isBlank(schemaDefinitionStr)) {
|
if (StringUtils.isBlank(definitionStr)) {
|
||||||
throw new RuntimeException("Schema definition must be provided for \"Avro\" schema type!");
|
throw new RuntimeException("Schema definition must be provided for \"Avro\" schema type!");
|
||||||
} else if (schemaDefinitionStr.startsWith(filePrefix)) {
|
|
||||||
try {
|
|
||||||
Path filePath = Paths.get(URI.create(schemaDefinitionStr));
|
|
||||||
schemaDefinitionStr = Files.readString(filePath, StandardCharsets.US_ASCII);
|
|
||||||
} catch (IOException ioe) {
|
|
||||||
throw new RuntimeException("Error reading the specified \"Avro\" schema definition file: " + definitionStr + ": " + ioe.getMessage());
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
return AVRO_SCHEMA_CACHE.computeIfAbsent(definitionStr, __ -> {
|
||||||
schema = AvroUtil.GetSchema_PulsarAvro("NBAvro", schemaDefinitionStr);
|
String schemaDefinitionStr = definitionStr;
|
||||||
|
if (schemaDefinitionStr.startsWith("file://")) {
|
||||||
|
try {
|
||||||
|
Path filePath = Paths.get(URI.create(schemaDefinitionStr));
|
||||||
|
schemaDefinitionStr = Files.readString(filePath, StandardCharsets.US_ASCII);
|
||||||
|
} catch (IOException ioe) {
|
||||||
|
throw new RuntimeException("Error reading the specified \"Avro\" schema definition file: " + definitionStr + ": " + ioe.getMessage());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return AvroUtil.GetSchema_PulsarAvro("NBAvro", schemaDefinitionStr);
|
||||||
|
});
|
||||||
} else {
|
} else {
|
||||||
throw new RuntimeException("Trying to create a \"Avro\" schema for a non-Avro schema type string: " + typeStr);
|
throw new RuntimeException("Trying to create a \"Avro\" schema for a non-Avro schema type string: " + typeStr);
|
||||||
}
|
}
|
||||||
|
|
||||||
return schema;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
///////
|
///////
|
||||||
|
Loading…
Reference in New Issue
Block a user