Flink 查询
Flink 查询
Iceberg 支持通过 Apache Flink 的 DataStream API 和 Table API 进行流式和批量读取。
使用 SQL 读取
Iceberg 在 Flink 中同时支持流式读取和批量读取。执行以下 SQL 命令可以在 streaming(流式)模式和 batch(批量)模式之间切换:
-- Execute the flink job in streaming mode for current session context
SET execution.runtime-mode = streaming;
-- Execute the flink job in batch mode for current session context
SET execution.runtime-mode = batch;Flink 批量读取
使用以下语句提交 Flink 批处理作业:
-- Execute the flink job in batch mode for current session context
SET execution.runtime-mode = batch;
SELECT * FROM sample;Flink 流式读取
Iceberg 支持在 Flink 流式作业中处理增量数据,作业可以从某个历史 snapshot-id 开始:
-- Submit the flink job in streaming mode for current session.
SET execution.runtime-mode = streaming;
-- Enable this switch because streaming read SQL will provide few job options in flink SQL hint options.
SET table.dynamic-table-options.enabled=true;
-- Read all the records from the iceberg current snapshot, and then read incremental data starting from that snapshot.
SELECT * FROM sample /*+ OPTIONS('streaming'='true', 'monitor-interval'='1s')*/ ;
-- Read all incremental data starting from the snapshot-id '3821550127947089987' (records from this snapshot will be excluded).
SELECT * FROM sample /*+ OPTIONS('streaming'='true', 'monitor-interval'='1s', 'start-snapshot-id'='3821550127947089987')*/ ;流式作业的 Flink SQL 中可以在 SQL hint 选项中设置一些参数,详情请参阅读取选项。
面向 SQL 的 FLIP-27 Source
以下是选择启用或禁用 FLIP-27 source 的 SQL 配置。
-- Opt out the FLIP-27 source.
-- Default is false for Flink 1.19 and below, and true for Flink 1.20 and above.
SET table.exec.iceberg.use-flip27-source = false;上述记录的所有其他 SQL 设置和选项均适用于 FLIP-27 source。
通过 SQL 读取分支和标签
可以通过指定选项来通过 SQL 读取分支和标签。有关更多详情,请参阅 Flink 配置。
--- Read from branch b1
SELECT * FROM table /*+ OPTIONS('branch'='b1') */ ;
--- Read from tag t1
SELECT * FROM table /*+ OPTIONS('tag'='t1') */;
--- Incremental scan from tag t1 to tag t2
SELECT * FROM table /*+ OPTIONS('streaming'='true', 'monitor-interval'='1s', 'start-tag'='t1', 'end-tag'='t2') */;使用 DataStream 读取
Iceberg 现已支持通过 Java API 进行流式或批量读取。
批量读取
此示例将在 Flink 批处理作业中从 Iceberg 表读取所有记录,然后将其打印到标准输出控制台:
StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://nn:8020/warehouse/path");
DataStream<RowData> batch = FlinkSource.forRowData()
.env(env)
.tableLoader(tableLoader)
.streaming(false)
.build();
// Print all records to stdout.
batch.print();
// Submit and execute this batch read job.
env.execute("Test Iceberg Batch Read");流式读取
本示例将读取从快照 ID '3821550127947089987' 开始的增量记录,并在 Flink 流式作业中打印到标准输出控制台:
StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://nn:8020/warehouse/path");
DataStream<RowData> stream = FlinkSource.forRowData()
.env(env)
.tableLoader(tableLoader)
.streaming(true)
.startSnapshotId(3821550127947089987L)
.build();
// Print all records to stdout.
stream.print();
// Submit and execute this streaming read job.
env.execute("Test Iceberg Streaming Read");还有其他可以设置的选项,请参阅 FlinkSource#Builder。
使用 DataStream 读取(FLIP-27 source)
FLIP-27 source 接口是在 Flink 1.12 中引入的。它旨在解决旧的 SourceFunction 流式 source 接口的若干不足,并统一了批处理与流处理执行的 source 接口。Flink 仓库中的大多数 source 连接器(如 Kafka、文件)都已迁移到 FLIP-27 接口。Flink 计划在不久的将来弃用旧的 SourceFunction 接口。
iceberg-flink 模块中新增了基于 FLIP-27 的 Flink IcebergSource。目前 FLIP-27 的 IcebergSource 仍为实验性功能。
批量读取
此示例将从 Iceberg 表中读取所有记录,然后在 Flink 批作业中打印到标准输出控制台:
StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://nn:8020/warehouse/path");
IcebergSource<RowData> source = IcebergSource.forRowData()
.tableLoader(tableLoader)
.assignerFactory(new SimpleSplitAssignerFactory())
.build();
DataStream<RowData> batch = env.fromSource(
source,
WatermarkStrategy.noWatermarks(),
"My Iceberg Source",
TypeInformation.of(RowData.class));
// Print all records to stdout.
batch.print();
// Submit and execute this batch read job.
env.execute("Test Iceberg Batch Read");流式读取
本示例将从表的最新快照(包含该快照)开始进行流式读取。每 60 秒轮询一次 Iceberg 表,以发现新的仅追加快照。目前尚不支持 CDC 读取。
StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://nn:8020/warehouse/path");
IcebergSource source = IcebergSource.forRowData()
.tableLoader(tableLoader)
.assignerFactory(new SimpleSplitAssignerFactory())
.streaming(true)
.streamingStartingStrategy(StreamingStartingStrategy.INCREMENTAL_FROM_LATEST_SNAPSHOT)
.monitorInterval(Duration.ofSeconds(60))
.build();
DataStream<RowData> stream = env.fromSource(
source,
WatermarkStrategy.noWatermarks(),
"My Iceberg Source",
TypeInformation.of(RowData.class));
// Print all records to stdout.
stream.print();
// Submit and execute this streaming read job.
env.execute("Test Iceberg Streaming Read");还有其他可通过 Java API 设置的选项,请参阅 IcebergSource#Builder。
通过 DataStream 读取分支和标签
分支和标签也可以通过 DataStream API 进行读取。
StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://nn:8020/warehouse/path");
// Read from branch
DataStream<RowData> batch = FlinkSource.forRowData()
.env(env)
.tableLoader(tableLoader)
.branch("test-branch")
.streaming(false)
.build();
// Read from tag
DataStream<RowData> batch = FlinkSource.forRowData()
.env(env)
.tableLoader(tableLoader)
.tag("test-tag")
.streaming(false)
.build();
// Streaming read from start-tag
DataStream<RowData> batch = FlinkSource.forRowData()
.env(env)
.tableLoader(tableLoader)
.streaming(true)
.startTag("test-tag")
.build();以 Avro GenericRecord 方式读取
FLIP-27 Iceberg source 提供了 AvroGenericRecordReaderFunction,用于将 Flink 的 RowData 转换为 Avro 的 GenericRecord。你可以使用该转换器,以 Avro GenericRecord DataStream 的形式从 Iceberg 表中读取数据。
请确保 flink-avro jar 已包含在 classpath 中。此外,iceberg-flink-runtime shaded 聚合 jar 无法使用,因为该 runtime jar 对 avro 包进行了 shade 处理。请改用未 shade 的 iceberg-flink jar。
TableLoader tableLoader = ...;
Table table;
try (TableLoader loader = tableLoader) {
loader.open();
table = loader.loadTable();
}
AvroGenericRecordReaderFunction readerFunction = AvroGenericRecordReaderFunction.fromTable(table);
IcebergSource<GenericRecord> source =
IcebergSource.<GenericRecord>builder()
.tableLoader(tableLoader)
.readerFunction(readerFunction)
.assignerFactory(new SimpleSplitAssignerFactory())
...
.build();
DataStream<Row> stream = env.fromSource(source, WatermarkStrategy.noWatermarks(),
"Iceberg Source as Avro GenericRecord", new GenericRecordAvroTypeInfo(avroSchema));发出水印
由数据源本身发出水印有多种好处,例如利用 Flink 水印对齐机制,或在并发读取多个数据文件时防止过早触发窗口计算。
通过设置 watermarkColumn 为 IcebergSource 启用水印生成。支持的列类型包括 timestamp、timestamptz 和 long。Iceberg 的 timestamp 或 timestamptz 本身就包含时间精度信息,因此无需指定时间单位。但 long 类型的列不包含时间单位信息,需使用 watermarkTimeUnit 来配置 long 列的转换方式。
水印基于数据文件存储的列指标生成,每个 split 发出一次。如果多个时间范围不同的小文件被合并到同一个 split 中,会增加乱序程度以及 Flink 状态中额外的数据缓冲。水印对齐的主要目的正是减少 Flink 状态中的乱序和多余的数据缓冲。因此,建议将 read.split.open-file-cost 设置为一个非常大的值,以避免将多个小文件合并到同一个 split 中。这样做的负面影响是(不合并小文件带来的)读取吞吐量下降,尤其是在存在大量小文件的情况下。在典型的状态化处理作业中,数据源的读取吞吐量通常不是瓶颈,因此这可能是一个合理的折中方案。
该功能依赖列级别的最小值/最大值统计信息。请确保在写入阶段为水印列生成统计数据。默认情况下,系统只为表的前 100 列收集列指标。如果水印列默认未启用统计数据,请在需要时使用以 write.metadata.metrics 开头的写入属性。
当水印用于窗口计算时,以下示例可能很有用。该数据源按顺序读取 Iceberg 数据文件,使用一个时间戳列并发出水印:
StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://nn:8020/warehouse/path");
DataStream<RowData> stream =
env.fromSource(
IcebergSource.forRowData()
.tableLoader(tableLoader)
// Watermark using timestamp column
.watermarkColumn("timestamp_column")
.build(),
// Watermarks are generated by the source, no need to generate it manually
WatermarkStrategy.<RowData>noWatermarks()
// Extract event timestamp from records
.withTimestampAssigner((record, eventTime) -> record.getTimestamp(pos, precision).getMillisecond()),
SOURCE_NAME,
TypeInformation.of(RowData.class));使用长整型事件时间列进行水位线对齐以读取 Iceberg 表的示例:
StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://nn:8020/warehouse/path");
DataStream<RowData> stream =
env.fromSource(
IcebergSource source = IcebergSource.forRowData()
.tableLoader(tableLoader)
// Disable combining multiple files to a single split
.set(FlinkReadOptions.SPLIT_FILE_OPEN_COST, String.valueOf(TableProperties.SPLIT_SIZE_DEFAULT))
// Watermark using long column
.watermarkColumn("long_column")
.watermarkTimeUnit(TimeUnit.MILLI_SCALE)
.build(),
// Watermarks are generated by the source, no need to generate it manually
WatermarkStrategy.<RowData>noWatermarks()
.withWatermarkAlignment(watermarkGroup, maxAllowedWatermarkDrift),
SOURCE_NAME,
TypeInformation.of(RowData.class));选项
读取选项
Flink 读取选项在配置 Flink IcebergSource 时传入:
IcebergSource.forRowData()
.tableLoader(TableLoader.fromCatalog(...))
.assignerFactory(new SimpleSplitAssignerFactory())
.streaming(true)
.streamingStartingStrategy(StreamingStartingStrategy.INCREMENTAL_FROM_LATEST_SNAPSHOT)
.startSnapshotId(3821550127947089987L)
.monitorInterval(Duration.ofMillis(10L)) // or .set("monitor-interval", "10s") \ set(FlinkReadOptions.MONITOR_INTERVAL, "10s")
.build()对于 Flink SQL,可以通过 SQL 提示(hint)传入读取选项,写法如下:
SELECT * FROM tableName /*+ OPTIONS('monitor-interval'='10s') */
...选项可以通过 Flink 配置传入,并应用于当前会话。请注意,并非所有选项都支持这种方式。
env.getConfig()
.getConfiguration()
.set(FlinkReadOptions.SPLIT_FILE_OPEN_COST_OPTION, 1000L);
...查看此处的所有选项:read-options
检查表
要检查表的历史记录、快照和其他元数据,Iceberg 支持元数据表。
元数据表的标识方式是在原始表名后附加元数据表名。例如,db.table 的历史记录通过 db.table$history 读取。
History(历史记录)
显示表的历史记录:
SELECT * FROM prod.db.table$history;| made_current_at | snapshot_id | parent_id | is_current_ancestor |
|---|---|---|---|
| 2019-02-08 03:29:51.215 | 5781947118336215154 | NULL | true |
| 2019-02-08 03:47:55.948 | 5179299526185056830 | 5781947118336215154 | true |
| 2019-02-09 16:24:30.13 | 296410040247533544 | 5179299526185056830 | false |
| 2019-02-09 16:32:47.336 | 2999875608062437330 | 5179299526185056830 | true |
| 2019-02-09 19:42:03.919 | 8924558786060583479 | 2999875608062437330 | true |
| 2019-02-09 19:49:16.343 | 6536733823181975045 | 8924558786060583479 | true |
Info
这里显示了一次被回滚的提交。 在本例中,快照 296410040247533544 和 2999875608062437330 拥有相同的父快照 5179299526185056830。快照 296410040247533544 已被回滚,因此不是当前表状态的祖先。
元数据日志条目
要显示表的元数据日志条目:
SELECT * from prod.db.table$metadata_log_entries;| timestamp | file | latest_snapshot_id | latest_schema_id | latest_sequence_number |
|---|---|---|---|---|
| 2022-07-28 10:43:52.93 | s3://.../table/metadata/00000-9441e604-b3c2-498a-a45a-6320e8ab9006.metadata.json | null | null | null |
| 2022-07-28 10:43:57.487 | s3://.../table/metadata/00001-f30823df-b745-4a0a-b293-7532e0c99986.metadata.json | 170260833677645300 | 0 | 1 |
| 2022-07-28 10:43:58.25 | s3://.../table/metadata/00002-2cc2837a-02dc-4687-acc1-b4d86ea486f4.metadata.json | 958906493976709774 | 0 | 2 |
快照
显示表的有效快照:
SELECT * FROM prod.db.table$snapshots;| committed_at | snapshot_id | parent_id | operation | manifest_list | summary |
|---|---|---|---|---|---|
| 2019-02-08 03:29:51.215 | 57897183625154 | null | append | s3://.../table/metadata/snap-57897183625154-1.avro | { added-records -> 2478404, total-records -> 2478404, added-data-files -> 438, total-data-files -> 438, flink.job-id -> 2e274eecb503d85369fb390e8956c813 } |
你也可以将快照与表历史进行关联。例如,下面的查询会显示表历史,并列出写入每个快照的应用 ID:
select
h.made_current_at,
s.operation,
h.snapshot_id,
h.is_current_ancestor,
s.summary['flink.job-id']
from prod.db.table$history h
join prod.db.table$snapshots s
on h.snapshot_id = s.snapshot_id
order by made_current_at;| made_current_at | operation | snapshot_id | is_current_ancestor | summary[flink.job-id] |
|---|---|---|---|---|
| 2019-02-08 03:29:51.215 | append | 57897183625154 | true | 2e274eecb503d85369fb390e8956c813 |
文件
显示表的当前数据文件:
SELECT * FROM prod.db.table$files;| content | file_path | file_format | spec_id | partition | record_count | file_size_in_bytes | column_sizes | value_counts | null_value_counts | nan_value_counts | lower_bounds | upper_bounds | key_metadata | split_offsets | equality_ids | sort_order_id |
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| 0 | s3:/.../table/data/00000-3-8d6d60e8-d427-4809-bcf0-f5d45a4aad96.parquet | PARQUET | 0 | {1999-01-01, 01} | 1 | 597 | [1 -> 90, 2 -> 62] | [1 -> 1, 2 -> 1] | [1 -> 0, 2 -> 0] | [] | [1 -> , 2 -> c] | [1 -> , 2 -> c] | null | [4] | null | null |
| 0 | s3:/.../table/data/00001-4-8d6d60e8-d427-4809-bcf0-f5d45a4aad96.parquet | PARQUET | 0 | {1999-01-01, 02} | 1 | 597 | [1 -> 90, 2 -> 62] | [1 -> 1, 2 -> 1] | [1 -> 0, 2 -> 0] | [] | [1 -> , 2 -> b] | [1 -> , 2 -> b] | null | [4] | null | null |
| 0 | s3:/.../table/data/00002-5-8d6d60e8-d427-4809-bcf0-f5d45a4aad96.parquet | PARQUET | 0 | {1999-01-01, 03} | 1 | 597 | [1 -> 90, 2 -> 62] | [1 -> 1, 2 -> 1] | [1 -> 0, 2 -> 0] | [] | [1 -> , 2 -> a] | [1 -> , 2 -> a] | null | [4] | null | null |
清单(Manifests)
显示表当前的文件清单:
SELECT * FROM prod.db.table$manifests;| path | length | partition_spec_id | added_snapshot_id | added_data_files_count | existing_data_files_count | deleted_data_files_count | partition_summaries |
|---|---|---|---|---|---|---|---|
| s3://.../table/metadata/45b5290b-ee61-4788-b324-b1e2735c0e10-m0.avro | 4479 | 0 | 6668963634911763636 | 8 | 0 | 0 | [[false,null,2019-05-13,2019-05-15]] |
注意:
manifests 表中
partition_summaries列内的各字段对应于清单列表中的field_summary结构体,顺序如下:contains_nullcontains_nanlower_boundupper_bound
contains_nan可能返回 null,表示文件的元数据中没有该信息。这通常发生在读取 V1 表时,因为 V1 表不会填充contains_nan。
分区
要显示表的当前分区:
SELECT * FROM prod.db.table$partitions;| partition | spec_id | record_count | file_count | total_data_file_size_in_bytes | position_delete_record_count | position_delete_file_count | equality_delete_record_count | equality_delete_file_count | last_updated_at(μs) | last_updated_snapshot_id |
|---|---|---|---|---|---|---|---|---|---|---|
| {20211001, 11} | 0 | 1 | 1 | 100 | 2 | 1 | 0 | 0 | 1633086034192000 | 9205185327307503337 |
| {20211002, 11} | 0 | 4 | 3 | 500 | 1 | 1 | 0 | 0 | 1633172537358000 | 867027598972211003 |
| {20211001, 10} | 0 | 7 | 4 | 700 | 0 | 0 | 0 | 0 | 1633082598716000 | 3280122546965981531 |
| {20211002, 10} | 0 | 3 | 2 | 400 | 0 | 0 | 1 | 1 | 1633169159489000 | 6941468797545315876 |
注意:对于未分区的表,partitions 表中将不包含 partition 和 spec_id 字段。
所有元数据表
这些表是当前快照对应的各元数据表的并集,返回所有快照范围内的元数据。
危险
"all" 元数据表可能会为同一个数据文件或清单文件返回多行,因为元数据文件可能同时属于多个表快照。
所有数据文件
展示表中所有的数据文件以及每个文件的元数据:
SELECT * FROM prod.db.table$all_data_files;所有 Manifest 文件
要显示表的所有 manifest 文件:
SELECT * FROM prod.db.table$all_manifests;| path | length | partition_spec_id | added_snapshot_id | added_data_files_count | existing_data_files_count | deleted_data_files_count | partition_summaries |
|---|---|---|---|---|---|---|---|
| s3://.../metadata/a85f78c5-3222-4b37-b7e4-faf944425d48-m0.avro | 6376 | 0 | 6272782676904868561 | 2 | 0 | 0 | [{false, false, 20210101, 20210101}] |
注意:
manifests 表中
partition_summaries列内的字段对应于清单列表中的field_summary结构体,顺序如下:contains_nullcontains_nanlower_boundupper_bound
contains_nan可能返回 null,表示文件的元数据中没有该信息。这种情况通常出现在读取 V1 表时,因为 V1 表不会填充contains_nan。
参考
显示表已知的快照引用(snapshot reference):
SELECT * FROM prod.db.table$refs;| name | type | snapshot_id | max_reference_age_in_ms | min_snapshots_to_keep | max_snapshot_age_in_ms |
|---|---|---|---|---|---|
| main | BRANCH | 4686954189838128572 | 10 | 20 | 30 |
| testTag | TAG | 4686954189838128572 | 10 | null | null |
评论
登录后参与评论
KnowForge