目录(Catalog)

Java 自定义目录

师成师成· 更新于 2026-09-28· 阅读 18 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

自定义目录(Catalog)

可以读取 Iceberg 表,其来源既可以是 HDFS 路径,也可以是 Hive 表;同时也可以用自定义元数据存储(metastore)来替代 Hive。具体步骤如下:

注意:要处理加密表,自定义目录必须满足一系列安全要求。

自定义表操作实现

继承 BaseMetastoreTableOperations,以提供读写元数据的实现。

示例:

class CustomTableOperations extends BaseMetastoreTableOperations {
  private String dbName;
  private String tableName;
  private Configuration conf;
  private FileIO fileIO;

  protected CustomTableOperations(Configuration conf, String dbName, String tableName) {
    this.conf = conf;
    this.dbName = dbName;
    this.tableName = tableName;
  }

  // The doRefresh method should provide implementation on how to get the metadata location
  @Override
  public void doRefresh() {

    // Example custom service which returns the metadata location given a dbName and tableName
    String metadataLocation = CustomService.getMetadataForTable(conf, dbName, tableName);

    // When updating from a metadata file location, call the helper method
    refreshFromMetadataLocation(metadataLocation);

  }

  // The doCommit method should provide implementation on how to update with metadata location atomically
  @Override
  public void doCommit(TableMetadata base, TableMetadata metadata) {
    String oldMetadataLocation = base.location();

    // Write new metadata using helper method
    String newMetadataLocation = writeNewMetadata(metadata, currentVersion() + 1);

    // Example custom service which updates the metadata location for the given db and table atomically
    CustomService.updateMetadataLocation(dbName, tableName, oldMetadataLocation, newMetadataLocation);

  }

  // The io method provides a FileIO which is used to read and write the table metadata files
  @Override
  public FileIO io() {
    if (fileIO == null) {
      fileIO = new HadoopFileIO(conf);
    }
    return fileIO;
  }
}

TableOperations 实例通常通过调用 Catalog.newTableOps(TableIdentifier) 获取。有关实现和加载自定义目录,请参阅下一节。

自定义目录实现

扩展 BaseMetastoreCatalog 以提供默认的仓库位置,并实例化 CustomTableOperations

示例:

public class CustomCatalog extends BaseMetastoreCatalog {

  private Configuration configuration;

  // must have a no-arg constructor to be dynamically loaded
  // initialize(String name, Map<String, String> properties) will be called to complete initialization
  public CustomCatalog() {
  }

  public CustomCatalog(Configuration configuration) {
    this.configuration = configuration;
  }

  @Override
  protected TableOperations newTableOps(TableIdentifier tableIdentifier) {
    String dbName = tableIdentifier.namespace().level(0);
    String tableName = tableIdentifier.name();
    // instantiate the CustomTableOperations
    return new CustomTableOperations(configuration, dbName, tableName);
  }

  @Override
  protected String defaultWarehouseLocation(TableIdentifier tableIdentifier) {

    // Can choose to use any other configuration name
    String tableLocation = configuration.get("custom.iceberg.warehouse.location");

    // Can be an s3 or hdfs path
    if (tableLocation == null) {
      throw new RuntimeException("custom.iceberg.warehouse.location configuration not set!");
    }

    return String.format(
            "%s/%s.db/%s", tableLocation,
            tableIdentifier.namespace().levels()[0],
            tableIdentifier.name());
  }

  @Override
  public boolean dropTable(TableIdentifier identifier, boolean purge) {
    // Example service to delete table
    CustomService.deleteTable(identifier.namespace().level(0), identifier.name());
  }

  @Override
  public void renameTable(TableIdentifier from, TableIdentifier to) {
    Preconditions.checkArgument(from.namespace().level(0).equals(to.namespace().level(0)),
            "Cannot move table between databases");
    // Example service to rename table
    CustomService.renameTable(from.namespace().level(0), from.name(), to.name());
  }

  // implement this method to read catalog name and properties during initialization
  public void initialize(String name, Map<String, String> properties) {
  }
}

在大多数计算引擎中,目录(Catalog)实现可以动态加载。对于 Spark 和 Flink,你可以指定 catalog-impl 目录属性来加载它。详情请阅读配置章节。对于 MapReduce,请实现 org.apache.iceberg.mr.CatalogLoader,并将 Hadoop 属性 iceberg.mr.catalog.loader.class 设置为该实现类来完成加载。如果你的目录必须读取 Hadoop 配置以访问特定的环境属性,请让该目录实现 org.apache.hadoop.conf.Configurable。

自定义 FileIO 实现

继承 FileIO 并提供数据文件的读写实现。

示例:

public class CustomFileIO implements FileIO {

  // must have a no-arg constructor to be dynamically loaded
  // initialize(Map<String, String> properties) will be called to complete initialization
  public CustomFileIO() {
  }

  @Override
  public InputFile newInputFile(String s) {
    // you also need to implement the InputFile interface for a custom input file
    return new CustomInputFile(s);
  }

  @Override
  public OutputFile newOutputFile(String s) {
    // you also need to implement the OutputFile interface for a custom output file
    return new CustomOutputFile(s);
  }

  @Override
  public void deleteFile(String path) {
    Path toDelete = new Path(path);
    FileSystem fs = Util.getFs(toDelete);
    try {
        fs.delete(toDelete, false /* not recursive */);
    } catch (IOException e) {
        throw new RuntimeIOException(e, "Failed to delete file: %s", path);
    }
  }

  // implement this method to read catalog properties during initialization
  public void initialize(Map<String, String> properties) {
  }
}

如果你已经在实现自己的 catalog,可以通过实现 TableOperations.io() 来使用自定义的 FileIO。此外,还可以通过指定 io-impl catalog 属性,在 HadoopCatalog 和 HiveCatalog 中动态加载自定义的 FileIO 实现。详情请参阅配置章节。如果你的 FileIO 必须读取 Hadoop 配置以访问某些环境属性,请让你的 FileIO 实现 org.apache.hadoop.conf.Configurable。

自定义位置提供器实现

继承 LocationProvider,并提供实现来确定写入数据的文件路径。

示例:

public class CustomLocationProvider implements LocationProvider {

  private String tableLocation;

  // must have a 2-arg constructor like this, or a no-arg constructor
  public CustomLocationProvider(String tableLocation, Map<String, String> properties) {
    this.tableLocation = tableLocation;
  }

  @Override
  public String newDataLocation(String filename) {
    // can use any custom method to generate a file path given a file name
    return String.format("%s/%s/%s", tableLocation, UUID.randomUUID().toString(), filename);
  }

  @Override
  public String newDataLocation(PartitionSpec spec, StructLike partitionData, String filename) {
    // can use any custom method to generate a file path given a partition info and file name
    return newDataLocation(filename);
  }
}

如果你已经在实现自己的 catalog,可以重写 TableOperations.locationProvider() 来使用你自定义的默认 LocationProvider。若要为特定表使用不同的自定义 location provider,请在创建表时通过表属性 write.location-provider.impl 指定实现类。

示例:

CREATE TABLE hive.default.my_table (
  id bigint,
  data string,
  category string)
USING iceberg
OPTIONS (
  'write.location-provider.impl'='com.my.CustomLocationProvider'
)
PARTITIONED BY (category);

自定义 IcebergSource

继承 IcebergSource 并提供实现,以便从 CustomCatalog 读取数据

示例:

public class CustomIcebergSource extends IcebergSource {

  @Override
  protected Table findTable(DataSourceOptions options, Configuration conf) {
    Optional<String> path = options.get("path");
    Preconditions.checkArgument(path.isPresent(), "Cannot open table: path is not set");

    // Read table from CustomCatalog
    CustomCatalog catalog = new CustomCatalog(conf);
    TableIdentifier tableIdentifier = TableIdentifier.parse(path.get());
    return catalog.loadTable(tableIdentifier);
  }
}

通过在其全限定名中更新 META-INF/services/org.apache.spark.sql.sources.DataSourceRegister 来注册 CustomIcebergSource

评论

登录后参与评论

正在加载评论…