You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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.jar
  • flink-connector-kafka_2.12-1.13.2.jar
  • flink-sql-connector-kafka_2.12-1.13.2.jar
  • kafka-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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.24 00:54:04