本文目录导读:

这是一个非常广泛的主题,为了给你提供有价值的案例,我将从架构设计、技术选型、核心代码和应用场景四个维度,构建一个基于 Apache Iceberg + Flink + Hive/Trino 的典型 Java 数据湖案例。
这个案例会模拟一个电商订单实时入湖的场景,并展示数据湖的核心能力(ACID、Time Travel、Schema Evolution)。
案例背景:电商实时订单数据湖
目标:将MySQL/App产生的实时订单数据,准实时地写入数据湖(Iceberg),并支持后续的OLAP分析(Trino)和历史回溯(Time Travel)。
技术栈:
- 存储层:Apache Iceberg(基于HDFS/S3)
- 计算层:Apache Flink(流式写入),Trino(OLAP查询)
- 元数据:Hive Metastore / AWS Glue
- 编程语言:Java 8+
架构设计图
[MySQL Binlog / Kafka]
|
|(CDC)
[Flink Job (Java)]
|
-----|-----
| Flink Iceberg Sink |
-----|-----
|
[Iceberg Table on HDFS/S3]
|
-----|-----
| Trino / Spark |
-----|-----
|
[BI / Ad-hoc Query / ML]
核心Java代码示例
1 Maven依赖 (pom.xml 关键部分)
<properties>
<iceberg.version>1.4.0</iceberg.version>
<flink.version>1.17.0</flink.version>
</properties>
<dependencies>
<!-- Flink Iceberg -->
<dependency>
<groupId>org.apache.iceberg</groupId>
<artifactId>iceberg-flink-runtime-1.17</artifactId>
<version>${iceberg.version}</version>
</dependency>
<!-- Iceberg Hive Catalog -->
<dependency>
<groupId>org.apache.iceberg</groupId>
<artifactId>iceberg-hive-metastore</artifactId>
<version>${iceberg.version}</version>
</dependency>
<!-- Flink Kafka Connector -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka</artifactId>
<version>${flink.version}</version>
</dependency>
<!-- Parquet & Avro (Iceberg 默认底层格式) -->
<dependency>
<groupId>org.apache.iceberg</groupId>
<artifactId>iceberg-parquet</artifactId>
<version>${iceberg.version}</version>
</dependency>
</dependencies>
2 Flink 流式写入 Iceberg (DataStream API)
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.data.RowData;
import org.apache.iceberg.flink.TableLoader;
import org.apache.iceberg.flink.sink.FlinkSink;
public class OrderStreamToIceberg {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(60000); // 每60秒做一次Checkpoint,保证Exactly-Once
// 1. 模拟从Kafka读取订单事件 (JSON -> RowData)
DataStream<OrderEvent> sourceStream = env
.addSource(new FlinkKafkaConsumer<>("orders_topic", new OrderDeserializationSchema(), kafkaProps))
.name("Kafka Order Source");
// 2. 将POJO转换为Flink内部RowData
DataStream<RowData> rowDataStream = sourceStream
.map(new OrderToRowDataMapper())
.name("Convert to RowData");
// 3. 配置Iceberg Table Loader
TableLoader tableLoader = TableLoader.fromHiveTable("hive_catalog", "ods_db", "orders_iceberg");
// 4. 使用Flink Sink写入Iceberg
FlinkSink.forRowData(rowDataStream)
.tableLoader(tableLoader)
.overwrite(false) // 追加式写入,不改历史
.distributionMode(DistributionMode.HASH) // 防止小文件
.writeParallelism(3)
.build();
env.execute("Real-time Order Ingestion to Iceberg");
}
}
3 Schema Evolution (Java 示例)
Iceberg 支持动态修改表结构,Flink 可以自动处理。
// 假设原始表有字段: id, user_id, amount, ts
// 某天新增了一个字段 "promotion_id"
import org.apache.iceberg.Table;
import org.apache.iceberg.hive.HiveCatalog;
import org.apache.iceberg.Schema;
import org.apache.iceberg.types.Types;
public class SchemaEvolutionDemo {
public static void main(String[] args) {
HiveCatalog catalog = new HiveCatalog();
catalog.setConf(hiveConf);
Table table = catalog.loadTable(TableIdentifier.of("ods_db", "orders_iceberg"));
// 添加新列 (Java 代码管理 Schema)
table.updateSchema()
.addColumn("promotion_id", Types.LongType.get())
.addColumn("delivery_note", Types.StringType.get())
.commit();
// 此时Flink写入新数据时,老数据行的promotion_id为null,新数据自动填充
System.out.println("Schema evolved successfully!");
}
}
4 Time Travel (Java API 查询历史快照)
import org.apache.iceberg.Table;
import org.apache.iceberg.DataFile;
import org.apache.iceberg.io.CloseableIterable;
import java.time.LocalDateTime;
import java.time.ZoneOffset;
// 查询2024年1月1日10点整的订单数据快照 (常用于回溯修复)
public class TimeTravelDemo {
public static void main(String[] args) {
Table table = ...; // 加载表
// 方式1: 根据时间戳查找快照
long targetTimestamp = LocalDateTime.of(2024, 1, 1, 10, 0)
.toInstant(ZoneOffset.UTC).toEpochMilli();
// 获取该时间点之前的最近一个快照ID
long snapshotId = table.snapshotAtTime(targetTimestamp).snapshotId();
// 方式2: 直接读取该快照的数据 (只包含已提交的文件)
CloseableIterable<FileScanTask> tasks = table.newScan()
.useSnapshot(snapshotId) // 关键: 指定快照ID
.planFiles();
// 读取文件内容 (省略具体读取逻辑)
System.out.println("查询时刻: " + targetTimestamp + ", 快照ID: " + snapshotId);
}
}
部署与运维命令 (Shell)
创建Hive Catalog对应的Iceberg数据库
# hive-sql
CREATE DATABASE IF NOT EXISTS ods_db;
CREATE TABLE ods_db.orders_iceberg (
order_id BIGINT,
user_id BIGINT,
amount DOUBLE,
ts TIMESTAMP
) STORED BY 'org.apache.iceberg.mr.hive.HiveIcebergStorageHandler'
TBLPROPERTIES('format-version'='2');
提交Flink任务
flink run -m yarn-cluster \ -c com.example.OrderStreamToIceberg \ -yjm 2048 -ytm 4096 \ /path/to/your-flink-iceberg-job.jar
使用Trino查询 Time Travel
-- 查询 2024-01-01 10:00:00 之前的最新数据 (默认) SELECT * FROM ods_db.orders_iceberg; -- 查询某时刻的数据 (Time Travel SQL) SELECT * FROM ods_db.orders_iceberg FOR SYSTEM_TIME AS OF '2024-01-01 10:00:00'; -- 查询变更历史 SELECT * FROM iceberg.ods_db."orders_iceberg$history";
关键设计与注意事项
| 关注点 | 实现方式 | 解决痛点 |
|---|---|---|
| ACID (并发写入) | Iceberg V2 格式 + Flink 两阶段提交 | 解决多个 Flink Job 同时写同一张表的脏读问题 |
| 小文件过多 | write.distribution-mode=hash + write.target-file-size-bytes=134217728(128MB) |
避免海量小文件拖垮NameNode |
| Schema兼容 | Iceberg自带Schema Evolution,列可增删改,不重写历史文件 | 业务字段频繁变更无需停机 |
| 数据回溯 | 快照隔离级别,默认读最新,可指定任意快照ID | 快速恢复误删数据/对比报表 |
| 压缩比 | 底层用Parquet + zstd | 存储成本降低60% |
实际业务效果
- 写入延迟:从Kafka到Iceberg可见,平均延迟 < 2分钟(受Checkpoint影响)。
- 查询延迟:Trino查询TB级数据,返回时间 < 5秒。
- 存储成本:相比HDFS原始文本,压缩比从1:3提升到1:8。
- 运维效率:Schema变更无需停服,历史数据自动兼容。
这个Java数据湖案例展示了:
- 实时入湖:Flink流式写入Iceberg,保证Exactly-Once语义。
- 查询解耦:Trino直接查询底层列存数据,无需ETL中间表。
- 核心特性:Time Travel、Schema Evolution、ACID事务在Java API中的实际调用方式。
如果你有具体的使用场景(比如CDC入湖、分区分桶优化),可以继续深入探讨。