JAVA中连接Flink写入paimon数据存至HDSF文件服务器上
·
文章目录
0. 需配置hadoop Configuration
0.1 配置无权限的hadoop configuration
/**
* 配置 hadoop 配置
*/
public static void checkHdfs(){
org.apache.hadoop.conf.Configuration hadoopConf = new org.apache.hadoop.conf.Configuration();
// 需与core-site.xml中的配置一致
hadoopConf.set("fs.defaultFS", "hdfs://ip:port");
}
0.2 配置kerberos认证的hadoop configuration
/**
* 配置 hadoop 关于 kerberos 的配置
*
* @param kdc5Conf kdc配置文件
* @param principal kerberos主体名称
* @param keyTab keytab密匙文件
*/
public static void checkHdfs(String kdc5Conf, String principal, String keyTab) {
try {
// 配置 Kerberos 系统属性
// 配置 kdc5Conf 配置文件
System.setProperty("java.security.krb5.conf", kdc5Conf);
// 配置开启 kerberos 调式日志
// System.setProperty("sun.security.krb5.debug", "true");
// 配置 hdfs
org.apache.hadoop.conf.Configuration hadoopConf = new org.apache.hadoop.conf.Configuration();
// 配置hadoop认证为kerberos
hadoopConf.set("hadoop.security.authentication", "kerberos");
// 配置hdfs的地址,需core-site.xml中的配置一致
hadoopConf.set("fs.defaultFS", "hdfs://ip:port");
// 配置hadoop中的namenode的kerberos主体
hadoopConf.set("dfs.namenode.kerberos.principal", principal);
// 配置hadoop中的datanode的kerberos主体
hadoopConf.set("dfs.datanode.kerberos.principal", principal);
// 重置并登录kerberos
UserGroupInformation.reset();
UserGroupInformation.setConfiguration(hadoopConf);
// 通过keytab登录kerberos
UserGroupInformation.loginUserFromKeytab(principal, keyTab);
// 获取当前用户
UserGroupInformation currentUser = UserGroupInformation.getCurrentUser();
// 对当前用户进行 Ticket 续期
currentUser.forceReloginFromKeytab();
System.out.println("kerberos认证成功。");
} catch (Exception e) {
System.out.println("kerberos认证失败。");
}
}
1. 创建TableEnvironment
/**
* 获取tableEnvironment
*
* @param kdc5Conf kdc配置文件
* @param principal kerberos主体名称
* @param keyTab keytab密匙文件
* @return TableEnvironment
*/
public static TableEnvironment getTableEnvironment(String kdc5Conf, String principal, String keyTab) {
Configuration flinkConf = new Configuration();
// 若需连接配置了kerberos的hadoop,需增加flink的配置, 若为配置kerberos认证,则下述四项均可不用配置
// 配置 kdc5conf 配置文件
flinkConf.setString(SecurityOptions.KERBEROS_KRB5_PATH, kdc5Conf);
// 配置 kerberos 主体(一般情况下配置namenode相同的principal)
flinkConf.setString(SecurityOptions.KERBEROS_LOGIN_PRINCIPAL, principal);
// 配置 keytab 文件
flinkConf.setString(SecurityOptions.KERBEROS_LOGIN_KEYTAB, keyTab);
// 配置不使用kerbeors缓存登录
flinkConf.setBoolean(SecurityOptions.KERBEROS_LOGIN_USETICKETCACHE, false);
// 构建TableEnvironment, 配置读取方式为 批式.inBatchMode()/流式.inStreamingMode()
EnvironmentSettings.Builder settingsBuilder = EnvironmentSettings.newInstance().inBatchMode();
// 应用自定义配置
EnvironmentSettings settings = settingsBuilder
.withConfiguration(flinkConf)
.build();
// 创建TableEnvironment
return TableEnvironment.create(settings);
}
2. 初始化 Paimon Catalog 并设置权限系统
/**
* 初始化 Paimon Catalog 并设置权限
* 若不需要配置权限,则无需执行Step 2之后的步骤
*
* @param tEnv TableEnvironment
* @param catalogName Catalog名称
* @param warehousePath 数据库存储路径
* @param initUser 待配置的用户名称
* @param initPassword 待配置的密码
*/
public static void initPaimonCatalogWithPrivilege(TableEnvironment tEnv, String catalogName, String warehousePath, String initUser, String initPassword) {
try {
String createCatalogSQL = String.format("CREATE CATALOG %s WITH (" +
" 'type' = 'paimon', 'warehouse' = '%s')",
catalogName, warehousePath);
// 执行语句
tEnv.executeSql(createCatalogSQL);
System.out.println("Step 1: 创建初始无权限的 Catalog 成功");
tEnv.useCatalog(catalogName);
System.out.println("Step 2: 切换到 Catalog: " + catalogName);
String initPrivilegeSQL = String.format( "CALL sys.init_file_based_privilege('%s')", initUser);
tEnv.executeSql(initPrivilegeSQL);
System.out.println("Step 3: 初始化权限系统成功,用户: " + initUser);
tEnv.useCatalog("default_catalog");
System.out.println("Step 4: 切换到默认 Catalog");
String dropCatalogSQL = String.format("DROP CATALOG %s", catalogName);
tEnv.executeSql(dropCatalogSQL);
System.out.println("Step 5: 删除无权限的 Catalog 成功");
String recreateCatalogSQL = String.format(
"CREATE CATALOG %s WITH (" +
" 'type' = 'paimon'," +
" 'warehouse' = '%s'," +
" 'user' = '%s'," +
" 'password' = '%s'" +
")",
catalogName,
warehousePath,
initUser,
initPassword
);
tEnv.executeSql(recreateCatalogSQL);
System.out.println("Step 6: 重新创建带权限的 Catalog 成功");
tEnv.useCatalog(catalogName);
System.out.println("Step 7: 切换到带权限的 Catalog: " + catalogName);
} catch (Exception e) {
System.err.println("初始化 Catalog 失败: " + e.getMessage());
}
}
3. 创建不同权限的用户
/**
* 创建不同权限的用户,若不需要配置用户可跳过该步骤
*
* @param tEnv TableEnvironment
* @param username 用户名
* @param password 密码
*/
public static void createPrivilegedUser(TableEnvironment tEnv, String username, String password) {
try {
String createUserSQL = String.format(
"CALL sys.create_privileged_user('%s', '%s')",
username,
password
);
tEnv.executeSql(createUserSQL);
System.out.println("创建用户成功: " + username);
} catch (Exception e) {
System.err.println("创建用户失败: " + e.getMessage());
}
}
4. 创建数据库
/**
* 创建数据库
*
* @param tEnv TableEnvironment
* @param databaseName 数据库名称
*/
public static void createDatabase(TableEnvironment tEnv, String databaseName) {
try {
String createDBSQL = String.format(
"CREATE DATABASE IF NOT EXISTS %s",
databaseName
);
tEnv.executeSql(createDBSQL);
System.out.println("创建数据库成功: " + databaseName);
} catch (Exception e) {
System.err.println("创建数据库失败: " + e.getMessage());
}
}
5. 创建表
/**
* 创建表
*
* @param tEnv TableEnvironment
* @param databaseName 数据库名称
* @param tableName 表名
* @param schema 表结构
*/
public static void createTable(TableEnvironment tEnv, String databaseName, String tableName, String schema) {
try {
// 使用数据库
tEnv.useDatabase(databaseName);
String createTableSQL = String.format(
"CREATE TABLE IF NOT EXISTS %s (%s)",
tableName,
schema
);
tEnv.executeSql(createTableSQL);
System.out.println("创建表成功: " + databaseName + "." + tableName);
} catch (Exception e) {
System.err.println("创建表失败: " + e.getMessage());
}
}
6. 插入简单数据样例
/**
* 插入数据
*
* @param tEnv TableEnvironment
* @param tableName 数据库表名
* @param batchSize 批次
* @param startNum 开始数
* @param endNum 结束数
* @param name 参数name
*/
public static void insertData(TableEnvironment tEnv, String tableName, int batchSize, int startNum, int endNum, String name) {
List<String> batchValues = new ArrayList<>();
try {
for (int i = startNum; i <= endNum; i++) {
// 添加当前记录到批次
batchValues.add("(" + i + ", '" + name + i + "', " + "'123@qq.com', " + "TIMESTAMP '2026-01-01 10:30:00'" + ")");
// 每达到设定批次数或者是最后一条记录时执行插入
if (batchValues.size() == batchSize || i == endNum) {
String sql = "INSERT INTO " + tableName + " VALUES " +
String.join(",", batchValues);
tEnv.executeSql(sql).await();
System.out.println("执行插入测试数据成功,累计插入: " + (i - startNum+ 1) + " 条数据");
// 清空批次
batchValues.clear();
}
}
System.out.println("插入测试数据成功。");
} catch (Exception e) {
System.out.println("插入数据失败。");
}
}
7. 授予用户不同级别的权限
7.1 授予 Catalog 级别的权限
/**
* 授予 Catalog 权限
*
* @param tEnv TableEnvironment
* @param username 用户名
* @param privilege 权限值 SELECT/INSERT/……
*/
public static void grantCatalogPrivilege(TableEnvironment tEnv, String username, String privilege) {
try {
String grantSQL = String.format(
"CALL sys.grant_privilege_to_user('%s', '%s')",
username,
privilege
);
tEnv.executeSql(grantSQL);
System.out.println("授予 Catalog 权限成功: " + username + " -> " + privilege);
} catch (Exception e) {
System.err.println("授予 Catalog 权限失败: " + e.getMessage());
}
}
7.2 授予 database 级别的权限
/**
* 授予数据库级别的权限
*
* @param tEnv TableEnvironment
* @param username 用户名
* @param privilege 权限值SELECT/INSERT/……
* @param databaseName 数据库名
*/
public static void grantDatabasePrivilege(TableEnvironment tEnv, String username, String privilege, String databaseName) {
try {
String grantSQL = String.format(
"CALL sys.grant_privilege_to_user('%s', '%s', '%s')",
username,
privilege,
databaseName
);
tEnv.executeSql(grantSQL);
System.out.println("授予数据库权限成功: " +
username + " -> " + databaseName + " -> " + privilege);
} catch (Exception e) {
System.err.println("授予数据库权限失败: " + e.getMessage());
}
}
7.2 授予 table 级别的权限
/**
* 授予表级别的权限
*
* @param tEnv TableEnvironment
* @param username 用户名
* @param privilege 权限值SELECT/INSERT/……
* @param databaseName 数据库名
* @param tableName 数据库表名
*/
public static void grantTablePrivilege(TableEnvironment tEnv, String username, String privilege, String databaseName, String tableName) {
try {
String grantSQL = String.format(
"CALL sys.grant_privilege_to_user('%s', '%s', '%s', '%s')",
username,
privilege,
databaseName,
tableName
);
tEnv.executeSql(grantSQL);
System.out.println("授予表权限成功: " +
username + " -> " + databaseName + "." + tableName + " -> " + privilege);
} catch (Exception e) {
System.err.println("授予表权限失败: " + e.getMessage());
}
}
8. 用户权限测试
/**
* 验证用户权限
*
* @param kdc5Conf kdc配置文件
* @param principal kerberos主体名称
* @param keyTab keytab密匙文件
* @param catalogName catalog名称
* @param warehousePath warehouse路径
* @param username 用户名
* @param password 密码
* @param database 数据库
* @param table 表名
* @param operation 对应的操作值
*/
public static void verifyUserPrivilege(String kdc5Conf, String principal, String keyTab, String catalogName, String warehousePath,
String username, String password, String database, String table, String operation) {
try {
TableEnvironment tEnv = getTableEnvironment(kdc5Conf, principal, keyTab);
String recreateCatalogSQL = String.format(
"CREATE CATALOG %s WITH (" +
" 'type' = 'paimon'," +
" 'warehouse' = '%s'," +
" 'user' = '%s'," +
" 'password' = '%s'" +
")",
catalogName,
warehousePath,
username,
password
);
tEnv.executeSql(recreateCatalogSQL);
System.out.println("创建 " + username + " 权限的 Catalog 成功");
tEnv.useCatalog(catalogName);
System.out.println("切换 " + username + " 权限的 Catalog 成功");
// 尝试执行操作来验证权限
if (database != null && table != null) {
// 验证表操作权限
tEnv.useDatabase(database);
String testSQL;
switch (operation.toUpperCase()) {
case "SELECT":
testSQL = String.format("SELECT * FROM %s LIMIT 1", table);
break;
case "INSERT":
// 这里需要根据表结构创建合适的插入语句
testSQL = String.format(
"INSERT INTO %s VALUES (101, 'test', '123@qq.com', TIMESTAMP '2026-01-01 10:30:00')", table);
break;
default:
testSQL = String.format("DESCRIBE %s", table);
}
try {
tEnv.executeSql(testSQL);
System.out.println("权限验证成功: " + username + " 有 " + database + "." + table + " 的" + operation + " 权限");
} catch (Exception e) {
System.out.println("权限验证失败: " + username + " 无 " +operation + " 权限");
}
}
} catch (Exception e) {
System.err.println("权限验证异常: " + e.getMessage());
}
}
9. 删除测试用户
/**
* 删除用户
*
* @param tEnv TableEnvironment
* @param username 用户名
*/
public static void dropPrivilegedUser(TableEnvironment tEnv, String username) {
try {
String dropUserSQL = String.format(
"CALL sys.drop_privileged_user('%s')",
username
);
tEnv.executeSql(dropUserSQL);
System.out.println("删除用户成功: " + username);
} catch (Exception e) {
System.err.println("删除用户失败: " + e.getMessage());
}
}
10. 完整调用流程
/**
* 完整调用流程
*
* @param tEnv TableEnvironment
* @param catalogName catalog名称
* @param warehousePath warehouse路径
* @param kdc5Conf kdc配置文件
* @param principal kerberos主体名称
* @param keyTab keytab密匙文件
*/
public static void runCompleteExample(TableEnvironment tEnv, String catalogName, String warehousePath,
String kdc5Conf, String principal, String keyTab) {
System.out.println("**** 初始化 Catalog 和权限系统 ****");
initPaimonCatalogWithPrivilege(tEnv, catalogName, warehousePath, "root", "root");
System.out.println("**** 创建测试用户 ****");
createPrivilegedUser(tEnv, "test", "test");
createPrivilegedUser(tEnv, "test1", "test1");
createPrivilegedUser(tEnv, "test2", "test2");
System.out.println("**** 创建数据库和表,并插入数据 ****");
createDatabase(tEnv, "testdb1");
createDatabase(tEnv, "testdb2");
tEnv.useDatabase("testdb1");
System.out.println("切换数据库testdb1");
// 创建表并分两次插入100条数据
createTable(tEnv, "testdb1", "user1",
"id BIGINT, name STRING, email STRING, register_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED");
insertData(tEnv, "user1", 50, 1, 100, "Alice");
createTable(tEnv, "testdb1", "user2",
"id BIGINT, name STRING, email STRING, register_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED");
insertData(tEnv, "user2", 50, 1, 100, "Bob");
tEnv.useDatabase("testdb2");
createTable(tEnv, "testdb2", "user_table1",
"id BIGINT, user_name STRING, email STRING, register_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED");
insertData(tEnv, "user_table1", 50, 1, 100, "Coco");
System.out.println("**** 授予权限 ****");
// 给 test 授予整个 Catalog 的 SELECT 权限
grantCatalogPrivilege(tEnv, "test", "SELECT");
// 给 test1 授予 testdb1 数据库的 SELECT 权限
grantDatabasePrivilege(tEnv, "test1", "SELECT", "testdb1");
// 给 test2 授予 testdb2 数据库 的 SELECT 权限
grantDatabasePrivilege(tEnv, "test2", "SELECT", "testdb2");
// 给 test2 授予 testdb1 中 user1 表的 SELECT 权限
grantTablePrivilege(tEnv, "test2", "SELECT", "testdb1", "user1");
tEnv.useDatabase("testdb1");
System.out.println("切换数据库testdb1");
System.out.println("**** 权限验证结果 ****");
System.out.println("==== root -- testdb1 -- user1====");
verifyUserPrivilege(kdc5Conf, principal, keyTab, catalogName, warehousePath,"root", "root",
"testdb1", "user1", "SELECT");
System.out.println("==== root -- testdb1 -- user2====");
verifyUserPrivilege(kdc5Conf, principal, keyTab, catalogName, warehousePath,"root", "root",
"testdb1", "user2", "SELECT");
System.out.println("==== root -- testdb2 -- user_table1====");
verifyUserPrivilege(kdc5Conf, principal, keyTab, catalogName, warehousePath,"root", "root",
"testdb2", "user_table1", "SELECT");
System.out.println("==== test -- testdb1 -- user1====");
verifyUserPrivilege(kdc5Conf, principal, keyTab, catalogName, warehousePath,"test", "test",
"testdb1", "user1", "SELECT");
System.out.println("==== test -- testdb1 -- user2====");
verifyUserPrivilege(kdc5Conf, principal, keyTab, catalogName, warehousePath,"test", "test",
"testdb1", "user2", "SELECT");
System.out.println("==== test -- testdb2 -- user_table1====");
verifyUserPrivilege(kdc5Conf, principal, keyTab, catalogName, warehousePath,"test", "test",
"testdb2", "user_table1", "SELECT");
System.out.println("==== test1 -- testdb1 -- user1====");
verifyUserPrivilege(kdc5Conf, principal, keyTab, catalogName, warehousePath,"test1", "test1",
"testdb1", "user1", "SELECT");
System.out.println("==== test1 -- testdb1 -- user2====");
verifyUserPrivilege(kdc5Conf, principal, keyTab, catalogName, warehousePath,"test1", "test1",
"testdb1", "user2", "SELECT");
System.out.println("==== test1 -- testdb2 -- user_table1====");
verifyUserPrivilege(kdc5Conf, principal, keyTab, catalogName, warehousePath,"test1", "test1",
"testdb2", "user_table1", "SELECT");
System.out.println("==== test2 -- testdb1 -- user1====");
verifyUserPrivilege(kdc5Conf, principal, keyTab, catalogName, warehousePath,"test2", "test2",
"testdb1", "user1", "SELECT");
System.out.println("==== test2 -- testdb1 -- user2====");
verifyUserPrivilege(kdc5Conf, principal, keyTab, catalogName, warehousePath,"test2", "test2",
"testdb1", "user2", "SELECT");
System.out.println("==== test2 -- testdb2 -- user_table1====");
verifyUserPrivilege(kdc5Conf, principal, keyTab, catalogName, warehousePath,"test2", "test2",
"testdb2", "user_table1", "SELECT");
System.out.println("**** 删除测试用户 ****");
dropPrivilegedUser(tEnv, "test");
dropPrivilegedUser(tEnv, "test1");
dropPrivilegedUser(tEnv, "test2");
}
11. 完整调用流程入口
public static void main(String[] args) {
String catalogName = "paimonCatalog";
String warehousePath = "hdfs://ip:port/xxx";
// 仅kerbero认证开启时才需要这三个参数
String kdc5Conf = "自己的kdc5配置文件";
String principal = "自己的主体";
String keyTab = "自己的keytab文件";
checkHdfs(kdc5Conf, principal, keyTab);
TableEnvironment tEnv = getTableEnvironment(kdc5Conf, principal, keyTab);
runCompleteExample(tEnv, catalogName, warehousePath, kdc5Conf, principal, keyTab);
}
控制台输出结果如下:
kerberos认证成功。
**** 初始化 Catalog 和权限系统 ****
Step 1: 创建初始无权限的 Catalog 成功
Step 2: 切换到 Catalog: paimonCatalog
Step 3: 初始化权限系统成功,用户: root
Step 4: 切换到默认 Catalog
Step 5: 删除无权限的 Catalog 成功
Step 6: 重新创建带权限的 Catalog 成功
Step 7: 切换到带权限的 Catalog: paimonCatalog
**** 创建测试用户 ****
创建用户成功: test
创建用户成功: test1
创建用户成功: test2
**** 创建数据库和表,并插入数据 ****
创建数据库成功: testdb1
创建数据库成功: testdb2
切换数据库testdb1
创建表成功: testdb1.user1
执行插入测试数据成功,累计插入: 50 条数据
执行插入测试数据成功,累计插入: 100 条数据
插入测试数据成功。
创建表成功: testdb1.user2
执行插入测试数据成功,累计插入: 50 条数据
执行插入测试数据成功,累计插入: 100 条数据
插入测试数据成功。
创建表成功: testdb2.user_table1
执行插入测试数据成功,累计插入: 50 条数据
执行插入测试数据成功,累计插入: 100 条数据
插入测试数据成功。
**** 授予权限 ****
授予 Catalog 权限成功: test -> SELECT
授予数据库权限成功: test1 -> testdb1 -> SELECT
授予数据库权限成功: test2 -> testdb2 -> SELECT
授予表权限成功: test2 -> testdb1.user1 -> SELECT
切换数据库testdb1
**** 权限验证结果 ****
==== test1 ====
创建 test1 权限的 Catalog 成功
切换 test1 权限的 Catalog: paimonCatalog
权限验证成功: test1 有 testdb1.user1 的SELECT 权限
创建 test1 权限的 Catalog 成功
切换 test1 权限的 Catalog: paimonCatalog
权限验证成功: test1 有 testdb1.user2 的SELECT 权限
创建 test1 权限的 Catalog 成功
切换 test1 权限的 Catalog: paimonCatalog
权限验证失败: test1 无 SELECT 权限
==== test2 ====
创建 test2 权限的 Catalog 成功
切换 test2 权限的 Catalog: paimonCatalog
权限验证成功: test2 有 testdb1.user1 的SELECT 权限
创建 test2 权限的 Catalog 成功
切换 test2 权限的 Catalog: paimonCatalog
权限验证失败: test2 无 SELECT 权限
创建 test2 权限的 Catalog 成功
切换 test2 权限的 Catalog: paimonCatalog
权限验证成功: test2 有 testdb2.user_table1 的SELECT 权限
**** 删除测试用户 ****
删除用户成功: test
删除用户成功: test1
删除用户成功: test2
Process finished with exit code 0

关注:1. 代码中引入的Flink相关包版本为1.19.1,paimon版本为1.0.1,若在执行代码中无法识别SQL语句导致的出错,需查询所需版本支持的SQL语句。2. 代码中演示的为基于kerberos认证的HDFS文件系统,若不需要,可自行配置默认无认证的方式。
12. HDFS文件服务器上查看写入的paimon数据
通过连接文件服务器,可以查看到在对应文件夹下创建的数据库、数据表、默认数据库、以及添加权限的.sys文件。

更多推荐


所有评论(0)