in src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java [109:124]
public void open() throws CatalogException {
if (mqAdminExt == null) {
try {
mqAdminExt = new DefaultMQAdminExt();
mqAdminExt.setNamesrvAddr(namesrvAddr);
mqAdminExt.setLanguage(LanguageCode.JAVA);
mqAdminExt.start();
} catch (MQClientException e) {
throw new CatalogException(
"Failed to create RocketMQ admin using :" + namesrvAddr, e);
}
}
if (schemaRegistryClient == null) {
schemaRegistryClient = SchemaRegistryClientFactory.newClient(schemaRegistryUrl, null);
}
}