Flink通过sql-client连接Kafka报错类路径下未找到kafka CatalogFactory
问题描述
尝试将Kafka与Flink对接,通过sql-client.sh运行任务,无论如何调整.yaml配置和依赖库,始终抛出如下错误:
Exception in thread "main" org.apache.flink.table.client.SqlClientException: Unexpected exception. This is a bug. Please consider filing an issue. at org.apache.flink.table.client.SqlClient.startClient(SqlClient.java:201) at org.apache.flink.table.client.SqlClient.main(SqlClient.java:161) Caused by: org.apache.flink.table.api.ValidationException: Unable to create catalog 'myKafka'. Catalog options are: 'type'='kafka' at org.apache.flink.table.factories.FactoryUtil.createCatalog(FactoryUtil.java:270) at org.apache.flink.table.client.gateway.context.LegacyTableEnvironmentInitializer.createCatalog(LegacyTableEnvironmentInitializer.java:217) at org.apache.flink.table.client.gateway.context.LegacyTableEnvironmentInitializer.lambda$initializeCatalogs$1(LegacyTableEnvironmentInitializer.java:120) at java.util.HashMap.forEach(HashMap.java:1289) at org.apache.flink.table.client.gateway.context.LegacyTableEnvironmentInitializer.initializeCatalogs(LegacyTableEnvironmentInitializer.java:117) at org.apache.flink.table.client.gateway.context.LegacyTableEnvironmentInitializer.initializeSessionState(LegacyTableEnvironmentInitializer.java:105) at org.apache.flink.table.client.gateway.context.SessionContext.create(SessionContext.java:233) at org.apache.flink.table.client.gateway.local.LocalContextUtils.buildSessionContext(LocalContextUtils.java:100) at org.apache.flink.table.client.gateway.local.LocalExecutor.openSession(LocalExecutor.java:91) at org.apache.flink.table.client.SqlClient.start(SqlClient.java:88) at org.apache.flink.table.client.SqlClient.startClient(SqlClient.java:187) ... 1 more Caused by: org.apache.flink.table.api.ValidationException: Could not find any factory for identifier 'kafka' that implements 'org.apache.flink.table.factories.CatalogFactory' in the classpath. Available factory identifiers are: generic_in_memory at org.apache.flink.table.factories.FactoryUtil.discoverFactory(FactoryUtil.java:319) at org.apache.flink.table.factories.FactoryUtil.getCatalogFactory(FactoryUtil.java:455) at org.apache.flink.table.factories.FactoryUtil.createCatalog(FactoryUtil.java:251) ... 11 more
相关配置与环境说明
使用的sql-conf配置如下(已隐去bootstrap servers等敏感信息):
catalogs: - name: myKafka type: kafka
library目录下已引入如下依赖jar包:
flink-avro-confluent-registry-1.13.2.jarflink-connector-kafka_2.12-1.13.2.jarflink-sql-connector-kafka_2.12-1.13.2.jarkafka-clients-2.0.0-cdh6.1.1.jar
当前使用的Flink版本为1.13.2,Kafka版本为2.0.0-cdh6.1.1。
根因说明
Flink 1.13官方未提供type=kafka的Catalog实现,你引入的Kafka连接器仅支持创建Kafka类型的表,不支持作为Catalog工厂加载,所以在配置文件中直接声明Kafka类型的Catalog会触发类路径找不到对应实现的报错。
解决方案
修改sql-conf.yaml使用hive catalog,在SQL客户端内部创建Kafka表即可正常使用。修改后的sql-conf.yaml配置如下:
execution: type: streaming result-mode: table planner: blink current-database: default current-catalog: myhive catalogs: - name: myhive type: hive hive-version: 2.1.1-cdh6.0.1 hive-conf-dir: /etc/hive/conf deployment: m: yarn-cluster yqu: ABC_XYZ
启动sql-client.sh后,在客户端内传入Kafka连接、主题、序列化格式等必要参数创建Kafka表,即可正常对接读写Kafka数据。
内容的提问来源于stack exchange,提问作者eemamedo
相关产品推荐
相关产品推荐

