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

Apache Flink与Doris集成问题:无法创建并切换至Doris Catalog

问题背景

我们正在开展PoC测试,评估Apache Flink-1.18.1与Apache Doris-2.0.4的集成能力,已将所有依赖添加至POM文件,Doris集群运行正常,可通过http://10.0.2.15:8030访问FE。在本地IDEA运行Java Maven项目时,尝试创建JDBC类型的Doris Catalog(demo_catalog)时抛出IllegalArgumentException,无法切换至该Catalog完成后续数据写入操作。


1. 当前代码中创建Catalog、表及插入数据的方式是否正确?

由于未提供具体代码,只能基于常见错误场景判断:

  • Catalog参数配置错误:JDBC Catalog必须包含name、type(值为jdbc)、default-database、username、password、base-url核心参数。若base-url格式错误(比如误用FE的HTTP端口8030,JDBC连接需用FE的9030端口)、参数缺失或值不合法,会直接抛出IllegalArgumentException。
  • 驱动类未正确指定:Doris兼容MySQL协议,需使用MySQL JDBC驱动(com.mysql.cj.jdbc.Driver),若未显式指定或依赖缺失,也会引发参数异常。
  • Catalog切换时机错误:需确保Catalog注册到Flink环境后再执行切换操作,若在Catalog未完成注册前调用useCatalog(),也可能触发异常。

2. 如何从Flink创建并访问Doris Catalog/表?示例代码与指引

方式一:使用JDBC Catalog连接Doris

Java示例代码:

import org.apache.flink.table.api.EnvironmentSettings;
import org.apache.flink.table.api.TableEnvironment;
import org.apache.flink.table.catalog.Catalog;
import org.apache.flink.table.catalog.jdbc.JdbcCatalog;

public class FlinkDorisJdbcCatalogDemo {
    public static void main(String[] args) {
        // 初始化流模式TableEnvironment
        EnvironmentSettings settings = EnvironmentSettings.inStreamingMode();
        TableEnvironment tableEnv = TableEnvironment.create(settings);

        // JDBC Catalog核心参数
        String catalogName = "demo_catalog";
        String defaultDb = "test_db";
        String username = "root";
        String password = "";
        // 注意:JDBC连接用FE的9030端口
        String baseUrl = "jdbc:doris://10.0.2.15:9030/" + defaultDb;

        // 创建并注册JDBC Catalog
        Catalog dorisJdbcCatalog = new JdbcCatalog(catalogName, defaultDb, username, password, baseUrl);
        tableEnv.registerCatalog(catalogName, dorisJdbcCatalog);

        // 切换至目标Catalog和数据库
        tableEnv.useCatalog(catalogName);
        tableEnv.useDatabase(defaultDb);

        // 创建Doris表(若不存在)
        tableEnv.executeSql("CREATE TABLE IF NOT EXISTS user_info (" +
                "id INT," +
                "name STRING," +
                "age INT" +
                ") WITH (" +
                "'connector' = 'jdbc'," +
                "'url' = '" + baseUrl + "'," +
                "'table-name' = 'user_info'," +
                "'username' = '" + username + "'," +
                "'password' = '" + password + "'" +
                ")");

        // 插入测试数据
        tableEnv.executeSql("INSERT INTO user_info VALUES (1, 'Alice', 25), (2, 'Bob', 30)").await();
    }
}

方式二:使用Doris官方Catalog(推荐,功能更完善)

Doris提供专属Flink Catalog,支持元数据自动同步等原生能力:

import org.apache.flink.table.api.EnvironmentSettings;
import org.apache.flink.table.api.TableEnvironment;

public class FlinkDorisNativeCatalogDemo {
    public static void main(String[] args) {
        EnvironmentSettings settings = EnvironmentSettings.inStreamingMode();
        TableEnvironment tableEnv = TableEnvironment.create(settings);

        // 注册Doris官方Catalog
        tableEnv.executeSql("CREATE CATALOG demo_catalog WITH (" +
                "'type' = 'doris'," +
                "'fe-nodes' = '10.0.2.15:8030'," + // 这里用FE的HTTP端口8030
                "'username' = 'root'," +
                "'password' = ''," +
                "'default-database' = 'test_db'" +
                ")");

        // 切换Catalog和数据库
        tableEnv.useCatalog("demo_catalog");
        tableEnv.useDatabase("test_db");

        // 创建Doris表
        tableEnv.executeSql("CREATE TABLE IF NOT EXISTS user_info (" +
                "id INT," +
                "name STRING," +
                "age INT" +
                ") WITH (" +
                "'connector' = 'doris'," +
                "'fenodes' = '10.0.2.15:8030'," +
                "'table.identifier' = 'test_db.user_info'," +
                "'username' = 'root'," +
                "'password' = ''," +
                "'sink.batch.size' = '1000'," +
                "'sink.batch.interval' = '1000'" +
                ")");

        // 插入测试数据
        tableEnv.executeSql("INSERT INTO user_info VALUES (1, 'Alice', 25), (2, 'Bob', 30)").await();
    }
}

官方文档指引:

  • Apache Flink官方文档:查看「Connectors > Doris Connector」章节,了解Catalog配置与数据读写细节。
  • Apache Doris官方文档:查看「集成 > Flink 集成」章节,获取版本兼容性、参数说明及最佳实践。

3. 是否遗漏了必要的前置配置?

以下是容易遗漏的关键配置:

  • 端口混淆:JDBC连接需使用FE的9030端口(MySQL协议端口),http://10.0.2.15:8030是FE的Web UI端口,二者不可混用。
  • 依赖配置:
    • 若使用JDBC Catalog,需添加Flink JDBC依赖与MySQL JDBC驱动:
      <dependency>
          <groupId>org.apache.flink</groupId>
          <artifactId>flink-connector-jdbc</artifactId>
          <version>1.18.1</version>
      </dependency>
      <dependency>
          <groupId>mysql</groupId>
          <artifactId>mysql-connector-java</artifactId>
          <version>8.0.33</version>
      </dependency>
      
    • 若使用Doris官方Catalog,需添加Doris Flink Connector依赖:
      <dependency>
          <groupId>org.apache.doris</groupId>
          <artifactId>flink-doris-connector-1.18_2.12</artifactId>
          <version>2.0.4</version>
      </dependency>
      
  • 权限配置:确保Doris用户(如root)拥有目标数据库的CREATE TABLE、INSERT权限,可通过Doris客户端执行GRANT ALL ON test_db.* TO 'root'@'%';授权。
  • 网络配置:确保本地IDEA或Flink集群节点能访问Doris FE的对应端口,无防火墙或安全组拦截。

内容的提问来源于stack exchange,提问作者Vikramsinh Shinde

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 15:07:25