Apache Flink与Doris集成问题:无法创建并切换至Doris Catalog
Flink 1.18.1 与 Doris 2.0.4 集成问题排查与解决方案
问题背景
我们正在开展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>
- 若使用JDBC Catalog,需添加Flink JDBC依赖与MySQL JDBC驱动:
- 权限配置:确保Doris用户(如
root)拥有目标数据库的CREATE TABLE、INSERT权限,可通过Doris客户端执行GRANT ALL ON test_db.* TO 'root'@'%';授权。 - 网络配置:确保本地IDEA或Flink集群节点能访问Doris FE的对应端口,无防火墙或安全组拦截。
内容的提问来源于stack exchange,提问作者Vikramsinh Shinde
相关产品推荐
相关产品推荐

