
最近在研究 AI BI智能数据分析 的落地实践。敬请期待后续专题实战系列《从零手把手教你搭建 AI 驱动的 BI 系统》将覆盖 Text2SQL、多轮对话、语义层、权限治理、生产级部署全链路代码可落地、坑点全复盘。Flink 系列文章一、Flink 专栏Flink 专栏系统介绍某一知识点并辅以具体的示例进行说明。1、Flink 部署系列本部分介绍Flink的部署、配置相关基础内容。2、Flink基础系列本部分介绍Flink 的基础部分比如术语、架构、编程模型、编程指南、基本的datastream api用法、四大基石等内容。3、Flik Table API和SQL基础系列本部分介绍Flink Table Api和SQL的基本用法比如Table API和SQL创建库、表用法、查询、窗口函数、catalog等等内容。4、Flik Table API和SQL提高与应用系列本部分是table api 和sql的应用部分和实际的生产应用联系更为密切以及有一定开发难度的内容。5、Flink 监控系列本部分和实际的运维、监控工作相关。二、Flink 示例专栏Flink 示例专栏是 Flink 专栏的辅助说明一般不会介绍知识点的信息更多的是提供一个一个可以具体使用的示例。本专栏不再分目录通过链接即可看出介绍的内容。两专栏的所有文章入口点击Flink 系列文章汇总索引文章目录Flink 系列文章一、通过 Table API 和 SQL Client 操作 HiveCatalog1、注册 Catalog1、方式一java实现2、方式二yaml配置2、修改当前的 Catalog 和数据库1、java实现2、sql3、列出可用的 Catalog1、java实现2、sql4、列出可用的数据库1、java实现2、sql5、列出可用的表1、java实现2、sql本文以示例展示了sql 和 table api 操作hivecatalog。一、通过 Table API 和 SQL Client 操作 HiveCatalog1、注册 Catalog用户可以访问默认创建的内存 Catalog default_catalog这个 Catalog 默认拥有一个默认数据库 default_database。 用户也可以注册其他的 Catalog 到现有的 Flink 会话中。以下通过api 和 配置文件注册catalog及配置。1、方式一java实现publicclassTestCreateHiveTable{publicstaticfinalStringtableNamealan_hivecatalog_hivedb_testTable;publicstaticfinalStringhive_create_table_sqlCREATE TABLE tableName (\n id INT,\n name STRING,\n age INT) TBLPROPERTIES (\n sink.partition-commit.delay5 s,\n sink.partition-commit.triggerpartition-time,\n sink.partition-commit.policy.kindmetastore,success-file);/** * param args * throws DatabaseAlreadyExistException * throws CatalogException */publicstaticvoidmain(String[]args)throwsException{StreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();StreamTableEnvironmenttenvStreamTableEnvironment.create(env);StringhiveConfDir/usr/local/bigdata/apache-hive-3.1.2-bin/conf;Stringnamealan_hive;// default 数据库名称StringdefaultDatabasedefault;HiveCataloghiveCatalognewHiveCatalog(name,defaultDatabase,hiveConfDir);tenv.registerCatalog(alan_hive,hiveCatalog);tenv.useCatalog(alan_hive);StringnewDatabaseNamealan_hivecatalog_hivedb;tenv.useDatabase(newDatabaseName);// 创建表tenv.getConfig().setSqlDialect(SqlDialect.HIVE);tenv.executeSql(hive_create_table_sql);// 插入数据StringinsertSQLinsert into alan_hivecatalog_hivedb_testTable values (1,alan,18);tenv.executeSql(insertSQL);// 查询数据StringselectSQLselect * from alan_hivecatalog_hivedb_testTable;Tabletabletenv.sqlQuery(selectSQL);table.printSchema();DataStreamTuple2Boolean,Rowresulttenv.toRetractStream(table,Row.class);result.print();env.execute();}}2、方式二yaml配置# 定义 catalogscatalogs:-name:alan_hivecatalogtype:hiveproperty-version:1hive-conf-dir:/usr/local/bigdata/apache-hive-3.1.2-bin/conf# 须包含 hive-site.xml# 改变表程序基本的执行行为属性。execution:planner:blink# 可选 blink 默认或 oldtype:streaming# 必选执行模式为 batch 或 streamingresult-mode:table# 必选table 或 changelogmax-table-result-rows:1000000# 可选table 模式下可维护的最大行数默认为 1000000小于 1 则表示无限制time-characteristic:event-time# 可选 processing-time 或 event-time 默认parallelism:1# 可选Flink 的并行数量默认为 1periodic-watermarks-interval:200# 可选周期性 watermarks 的间隔时间默认 200 msmax-parallelism:16# 可选Flink 的最大并行数量默认 128min-idle-state-retention:0# 可选表程序的最小空闲状态时间max-idle-state-retention:0# 可选表程序的最大空闲状态时间current-catalog:alan_hivecatalog# 可选当前会话 catalog 的名称默认为 default_catalogcurrent-database:viewtest_db# 可选当前 catalog 的当前数据库名称默认为当前 catalog 的默认数据库restart-strategy:# 可选重启策略restart-strategytype:fallback# 默认情况下“回退”到全局重启策略2、修改当前的 Catalog 和数据库Flink 始终在当前的 Catalog 和数据库中寻找表、视图和 UDF。1、java实现代码片段只列出了关键的代码。StreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();StreamTableEnvironmenttenvStreamTableEnvironment.create(env);StringcatalogNamealan_hive;StringdefaultDatabasedefault;StringdatabaseNameviewtest_db;StringhiveConfDir/usr/local/bigdata/apache-hive-3.1.2-bin/conf;HiveCataloghiveCatalognewHiveCatalog(catalogName,defaultDatabase,hiveConfDir);tenv.registerCatalog(catalogName,hiveCatalog);tenv.useCatalog(catalogName);hiveCatalog.createDatabase(databaseName,newCatalogDatabaseImpl(newHashMap(),hiveConfDir){},true);// tenv.executeSql(create database databaseName);tenv.useDatabase(databaseName);2、sqlFlinkSQLUSECATALOG alan_hive;FlinkSQLUSEviewtest_db;通过提供全限定名 catalog.database.object 来访问不在当前 Catalog 中的元数据信息。javatenv.from(not_the_current_catalog.not_the_current_db.my_table);sqlFlinkSQLSELECT*FROMnot_the_current_catalog.not_the_current_db.my_table;3、列出可用的 Catalog1、java实现tenv.listCatalogs();2、sqlshowcatalogs;4、列出可用的数据库1、java实现tenv.listDatabases();2、sqlshowdatabases;5、列出可用的表1、java实现tenv.listTables();2、sqlshowtables;以上本文以示例展示了sql 和 table api 操作hivecatalog。