flink安装部署,写入paimon数据,Java读取数据
文章目录
1. flink安装部署和启动
1.1 下载flink安装包,并解压到文件目录下
根据需要下载flink安装包,链接: flink
tar -zxvf flink-1.19.1-bin-scala_2.12.tgz
1.2 进入解压后的flink安装目录下的bin目录,启动flink服务
cd bin
cd ./start-cluster.sh

2. flink-sql写入paimon数据
2.1 导入相关包至flink安装目录下的lib目录
1> 存储本地文件系统:若paimon数据存储于本地文件系统,则导入paimon-flink包,如:paimon-flink-1.19-1.0.1.jar包,版本依据实际需要导入。
2> 存储HDFS文件系统:若paimon数据存储于HDFS文件系统,则需导入hadoop相关的包,hadoop-hdfs-3.3.6.jar、hadoop-client-3.3.6.jar、hadoop-common-3.3.6.jar、hadoop-client-runtime-3.3.4.jar、hadoop-client-api-3.3.5.jar。
2.2 进入flink安装目录下的bin目录,启动flink的client客户端
./sql-client.sh

2.3. 创建paimon的catalog并使用
若存储在服务器文件系统上,将warehouse的值改为文件夹路径即可,如/home/paimondb/,此处写入的是HDFS文件系统。
create catalog paimonCatalog with('type'='paimon','warehouse'='hdfs://localhost:9000/paimon/warehouse');

创建成功后,当前会话中会包含两个catalog,一个是默认的default_catalog,另一个是自己创建成功的paimonCatalog。
2.4. 创建数据库并使用
CREATE DATABASE IF NOT EXISTS testdb;
USE testdb;

同catalog一样,新建的catalog下会生成一个默认的数据库default,通过【show databases;】命令可以查看创建的表,也可也通过系统文件或HDFS文件系统中查看数据库格式为:【/paimon/warehouse/default.db】和【/paimon/warehouse/testdb.db】。
2.5 创建简单表并插入3条数据
CREATE TABLE user_table (user_id BIGINT, user_name STRING, email STRING, register_time TIMESTAMP(3), idCard STRING, PRIMARY KEY (user_id) NOT ENFORCED);
INSERT INTO user_table VALUES (1, 'test', '123@qq.com', TIMESTAMP '2025-01-15 00:00:00', '440521199703011234');
INSERT INTO user_table VALUES (2, 'test2', '123@qq.com', TIMESTAMP '2026-01-15 00:00:00', '440521199703011234'),(3, '测试3', '123@qq.com', TIMESTAMP '2026-01-15 00:00:00', '440521199703011234');



默认流式展示查询结果,可通过如下配置变换查询结果的模式:
set 'sql-client.execution.result-mode' = 'tableau';
reset 'execution.checkpoionting.interval';
set 'execution.runtime-mode' = 'batch';
select * from user_table;

3. java代码读取paimon数据
public static void main(String[] args) {
String catalogName = "paimon_catalog";
String warehousePath = "hdfs://localhost:9000/paimon/warehouse";
org.apache.hadoop.conf.Configuration hadoopConf = new org.apache.hadoop.conf.Configuration();
hadoopConf.set("fs.defaultFS", "hdfs://localhost:9000");
Configuration flinkConf = new Configuration();
// 构建TableEnvironment, 配置读取方式为 批式.inBatchMode()/流式.inStreamingMode()
EnvironmentSettings.Builder settingsBuilder = EnvironmentSettings.newInstance().inBatchMode();
// 应用自定义配置
EnvironmentSettings settings = settingsBuilder
.withConfiguration(flinkConf)
.build();
// 创建TableEnvironment
TableEnvironment tEnv = TableEnvironment.create(settings);
System.out.println("1. tEnv 创建成功");
// 1. 创建初始 Catalog(无权限控制)
String createCatalogSQL = String.format("CREATE CATALOG %s WITH ('type' = 'paimon', 'warehouse' = '%s')", catalogName, warehousePath );
tEnv.executeSql(createCatalogSQL);
System.out.println("2. catalog 创建成功");
// 切换到创建的catalog下
tEnv.useCatalog(catalogName);
// 获取所有的数据库
String[] databaseList = tEnv.listDatabases();
for(String database : databaseList) {
// 一般不会使用默认数据库,所以可以跳过,若有需要可删除此跳过语句
if (database.equals("default")) {
continue;
}
// 切换数据库
tEnv.useDatabase(database);
System.out.println("3. use database: " + database);
String[] tableList = tEnv.listTables();
for (String table : tableList) {
System.out.println("================== table: " + table + " ==================");
String testSQL = String.format("SELECT * FROM %s ", table);
tEnv.executeSql(testSQL).print();
System.out.println("================== ==================== ==================");
}
System.out.println("4. use database: " + database + " 结束");
}
}

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。
更多推荐



所有评论(0)