Java 自定义目录
自定义目录(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
评论
登录后参与评论
KnowForge