我竟然在本地跑通了 Iceberg!Java 操作全流程曝光
我竟然在本地跑通了 Iceberg!Java 操作全流程曝光
什么是 Iceberg
Iceberg 是一种适用于大型分析表的高性能格式。Iceberg 为大数据带来了 SQL 表的可靠性和简洁性,同时使 Spark、Trino、Flink、Presto、Hive 和 Impala 等引擎能够同时安全地处理同一张表。
环境搭建
最快的入门方法是使用 docker-compose 文件,该文件使用了tabulario/spark-iceberg镜像,其中包含一个已配置 Iceberg 目录的本地 Spark 集群。要使用此功能,您需要安装Docker CLI以及Docker Compose CLI。
获得这些后,将下面的 yaml 保存到名为的文件中docker-compose.yml:
services:
spark-iceberg:
image: tabulario/spark-iceberg
container_name: spark-iceberg
build: spark/
networks:
iceberg_net:
depends_on:
- rest
- minio
volumes:
- ./warehouse:/home/iceberg/warehouse
- ./notebooks:/home/iceberg/notebooks/notebooks
environment:
- AWS_ACCESS_KEY_ID=admin
- AWS_SECRET_ACCESS_KEY=password
- AWS_REGION=us-east-1
ports:
- 8888:8888
- 8080:8080
- 10000:10000
- 10001:10001
rest:
image: apache/iceberg-rest-fixture
container_name: iceberg-rest
networks:
iceberg_net:
ports:
- 8181:8181
environment:
- AWS_ACCESS_KEY_ID=admin
- AWS_SECRET_ACCESS_KEY=password
- AWS_REGION=us-east-1
- CATALOG_WAREHOUSE=s3://warehouse/
- CATALOG_IO__IMPL=org.apache.iceberg.aws.s3.S3FileIO
- CATALOG_S3_ENDPOINT=http://minio:9000
minio:
image: minio/minio
container_name: minio
environment:
- MINIO_ROOT_USER=admin
- MINIO_ROOT_PASSWORD=password
- MINIO_DOMAIN=minio
networks:
iceberg_net:
aliases:
- warehouse.minio
ports:
- 9001:9001
- 9000:9000
command: ["server", "/data", "--console-address", ":9001"]
mc:
depends_on:
- minio
image: minio/mc
container_name: mc
networks:
iceberg_net:
environment:
- AWS_ACCESS_KEY_ID=admin
- AWS_SECRET_ACCESS_KEY=password
- AWS_REGION=us-east-1
entrypoint: |
/bin/sh -c "
until (/usr/bin/mc alias set minio http://minio:9000 admin password) do echo '...waiting...' && sleep 1; done;
/usr/bin/mc rm -r --force minio/warehouse;
/usr/bin/mc mb minio/warehouse;
/usr/bin/mc policy set public minio/warehouse;
tail -f /dev/null
"
networks:
iceberg_net:
接下来,使用以下命令启动 docker 容器:
docker-compose up
快速开始
新建一个java项目,并引入依赖。
添加依赖
<dependency>
<groupId>org.apache.iceberg</groupId>
<artifactId>iceberg-api</artifactId>
<version>1.3.1</version>
</dependency>
<dependency>
<groupId>org.apache.iceberg</groupId>
<artifactId>iceberg-core</artifactId>
<version>1.3.1</version>
</dependency>
<dependency>
<groupId>org.apache.iceberg</groupId>
<artifactId>iceberg-aws</artifactId>
<version>1.3.1</version>
</dependency>
<dependency>
<groupId>org.apache.iceberg</groupId>
<artifactId>iceberg-data</artifactId>
<version>1.3.1</version>
</dependency>
<dependency>
<groupId>org.apache.iceberg</groupId>
<artifactId>iceberg-parquet</artifactId>
<version>1.3.1</version>
</dependency>
<dependency>
<groupId>org.apache.httpcomponents.client5</groupId>
<artifactId>httpclient5</artifactId>
<version>5.3</version>
</dependency>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-client</artifactId>
<version>3.3.6</version>
</dependency>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-aws</artifactId>
<version>3.4.0</version>
</dependency>
catalog 配置项
要加载catalog,首先必须构建一个Map来配置它。所需的属性取决于使用的catalog类型。我们使用的是 REST catalog ,所有 catalog 通常都需要两个属性:仓库位置和 FileIO 实现。我们将使用与 S3 兼容的 Minio 容器。
下面生成 catalog 属性 Map 来配置我们的 RestCatalog。
import org.apache.iceberg.catalog.Catalog;
import org.apache.hadoop.conf.Configuration;
import org.apache.iceberg.CatalogProperties;
import org.apache.iceberg.rest.RESTCatalog;
import org.apache.iceberg.aws.AwsProperties;
Map<String, String> properties = new HashMap<>();
properties.put(CatalogProperties.CATALOG_IMPL, "org.apache.iceberg.rest.RESTCatalog");
properties.put(CatalogProperties.URI, "http://localhost:8181");
properties.put(CatalogProperties.WAREHOUSE_LOCATION, "s3a://warehouse");
properties.put(CatalogProperties.FILE_IO_IMPL, "org.apache.iceberg.aws.s3.S3FileIO");
properties.put(AwsProperties.S3FILEIO_ENDPOINT, "http://localhost:9000");
properties.put(AwsProperties.S3FILEIO_PATH_STYLE_ACCESS, "true");
properties.put(AwsProperties.CLIENT_REGION, "us-east-1");
properties.put(AwsProperties.S3FILEIO_ACCESS_KEY_ID, "admin");
properties.put(AwsProperties.S3FILEIO_SECRET_ACCESS_KEY, "password");
接下来,初始化 catalog ,为其设置名称并将包含我们配置的属性Map传递给它。
RESTCatalog catalog = new RESTCatalog();
Configuration conf = new Configuration();
catalog.setConf(conf);
catalog.initialize("demo", properties);
这样我们现在有一个 catalog 实例,其中包含列出、创建、重命名和删除表等操作。
定义 Schema 和 Partition Spec
我们将创建一个表,但首先,我们必须定义表的schema。让我们创建一个包含5列的简单schema (event_id 、 username、userid、 api_version、command)。
import org.apache.iceberg.Schema;
import org.apache.iceberg.types.Types;
Schema schema = new Schema(
Types.NestedField.optional(1, "event_id", Types.StringType.get()),
Types.NestedField.optional(2, "username", Types.StringType.get()),
Types.NestedField.optional(3, "userid", Types.IntegerType.get()),
Types.NestedField.optional(4, "api_version", Types.StringType.get()),
Types.NestedField.optional(5, "command", Types.StringType.get())
);
此外,让我们构建一个Partition Spec,定义列上的根据 api_version 来分区 。
import org.apache.iceberg.PartitionSpec;
PartitionSpec spec = PartitionSpec.builderFor(schema)
.hour("api_version")
.build();
创建表
使用已经创建的 Schema 和 Partition Spec,我们现在可以创建表了。我们将创建一个“webapp”命名空间,并在该命名空间中创建表user_events。
import org.apache.iceberg.catalog.Namespace;
import org.apache.iceberg.catalog.TableIdentifier;
Namespace namespace = Namespace.of("webapp");
TableIdentifier name = TableIdentifier.of(namespace, "user_events");
现在,让我们创建我们的表!
catalog.createTable(name, schema, spec)
如果我们在目录上调用 listTables 方法,我们可以在列表中看到我们新创建的表。
List<TableIdentifier> tables = catalog.listTables(namespace);
System.out.println(tables)
创建记录
Iceberg 数据模块有一个 GenericRecord 类,可以轻松创建与 Iceberg 模式匹配的记录。让我们使用该类创建一个包含一些记录的列表。
import java.util.UUID;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableMap;
import org.apache.iceberg.data.GenericRecord;
GenericRecord record = GenericRecord.create(schema);
ImmutableList.Builder<GenericRecord> builder = ImmutableList.builder();
builder.add(record.copy(ImmutableMap.of("event_id", UUID.randomUUID().toString(), "username", "Bruce", "userid", 1, "api_version", "1.0", "command", "grapple")));
builder.add(record.copy(ImmutableMap.of("event_id", UUID.randomUUID().toString(), "username", "Wayne", "userid", 1, "api_version", "1.0", "command", "glide")));
builder.add(record.copy(ImmutableMap.of("event_id", UUID.randomUUID().toString(), "username", "Clark", "userid", 1, "api_version", "2.0", "command", "fly")));
builder.add(record.copy(ImmutableMap.of("event_id", UUID.randomUUID().toString(), "username", "Kent", "userid", 1, "api_version", "1.0", "command", "land")));
ImmutableList<GenericRecord> records = builder.build();
现在我们有了记录列表,我们准备将它们写入文件。
import org.apache.iceberg.Files;
import org.apache.iceberg.io.DataWriter;
import org.apache.iceberg.io.OutputFile;
import org.apache.iceberg.parquet.Parquet;
import org.apache.iceberg.data.parquet.GenericParquetWriter;
String filepath = table.location() + "/" + UUID.randomUUID().toString();
OutputFile file = table.io().newOutputFile(filepath);
DataWriter<GenericRecord> dataWriter =
Parquet.writeData(file)
.schema(schema)
.createWriterFunc(GenericParquetWriter::buildWriter)
.overwrite()
.withSpec(PartitionSpec.unpartitioned())
.build();
try {
for (GenericRecord record : builder.build()) {
dataWriter.write(record);
}
} finally {
dataWriter.close();
}
dataWriter关闭后,即可将其转换为数据文件。数据文件包含 Iceberg 表所需的所有元数据,例如路径、
大小、记录数以及列级统计信息。
import org.apache.iceberg.DataFile;
DataFile dataFile = dataWriter.toDataFile();
然后可以将数据文件直接附加到 Iceberg 表并提交。
import org.apache.iceberg.catalog.Namespace;
import org.apache.iceberg.catalog.TableIdentifier;
import org.apache.iceberg.Table;
Namespace webapp = Namespace.of("webapp");
TableIdentifier name = TableIdentifier.of(webapp, "user_events");
Table tbl = catalog.loadTable(name);
tbl.newAppend().appendFile(dataFile).commit()
查询记录
可以使用表扫描来读取表以确认数据已提交。
import org.apache.iceberg.io.CloseableIterable;
import org.apache.iceberg.data.Record;
import org.apache.iceberg.data.IcebergGenerics;
CloseableIterable<Record> result = IcebergGenerics.read(tbl).build();
for (Record r: result) {
System.out.println(r);
}
通常情况下,全表扫描并非最佳选择,尤其是对于包含海量数据的表。使用 Iceberg 表达式,您可以通过将 where
子句链接到扫描构建器来为表扫描添加过滤器。让我们添加一个表达式,将扫描结果过滤为仅包含“userid==1”的记录。
CloseableIterable<Record> result1 = IcebergGenerics.read(table)
.where(Expressions.equal("userid", 1))
.build();
for (Record r: result1) {
System.out.println(r);
}
查询文件
要将 Iceberg 表中的数据提取到内存中,使用 IcebergGenerics 非常有效。但是,当将 Iceberg 集成到计算框架或查询引擎时,生成与
特定表达式匹配的文件列表更为有用。为此,您可以构造一个 TableScan 生成一组任务的对象。
import org.apache.iceberg.CombinedScanTask;
import org.apache.iceberg.TableScan;
TableScan scan = table.newScan();
就像 IcebergGenerics 一样 ,您可以将表达式应用于扫描。让我们添加一个表达式,过滤扫描结果,使其仅包含“userid==1”记录。
import org.apache.iceberg.expressions.Expressions;
TableScan filteredScan = scan.filter(Expressions.equal("userid", 1)).select("message")
现在我们可以从过滤后的扫描中检索任务列表。
Iterable<CombinedScanTask> result = filteredScan.planTasks();
CombinedScanTask 包含 FileScanTask 不同类型文件的实例,例如数据文件。让我们提取第一个文件并进行检查。
import org.apache.iceberg.DataFile;
CombinedScanTask task = result.iterator().next();
DataFile dataFile = task.files().iterator().next().file();
System.out.println(dataFile);
该 DataFile 实例包含有关文件的大量信息,例如:
– 路径
– 格式
– 大小
– 所属分区
– 包含的记录数
总结
本文介绍了 Iceberg 的核心特性,同时,通过一个完整的 Java 示例,演示了如何搭建 Iceberg 的本地开发环境,配置 REST Catalog,定义表结构,写入和读取数据,并执行高效的表扫描操作。
Iceberg 不仅解决了传统数据湖中常见的性能瓶颈和一致性问题,还为构建现代数据架构提供了坚实的基础。对于希望在大数据平台上实现高性能、高并发、可扩展的数据管理方案的开发者和架构师来说,Iceberg 是一个非常值得深入研究和采用的工具。
更多推荐


所有评论(0)