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文件。
HDFS文件服务器查看paimon数据

更多推荐