Apache Flink

Flink 查询

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

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

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_atsnapshot_idparent_idis_current_ancestor
2019-02-08 03:29:51.2155781947118336215154NULLtrue
2019-02-08 03:47:55.94851792995261850568305781947118336215154true
2019-02-09 16:24:30.132964100402475335445179299526185056830false
2019-02-09 16:32:47.33629998756080624373305179299526185056830true
2019-02-09 19:42:03.91989245587860605834792999875608062437330true
2019-02-09 19:49:16.34365367338231819750458924558786060583479true

Info

这里显示了一次被回滚的提交。 在本例中,快照 296410040247533544 和 2999875608062437330 拥有相同的父快照 5179299526185056830。快照 296410040247533544 已被回滚,因此不是当前表状态的祖先。

元数据日志条目

要显示表的元数据日志条目:

SELECT * from prod.db.table$metadata_log_entries;
timestampfilelatest_snapshot_idlatest_schema_idlatest_sequence_number
2022-07-28 10:43:52.93s3://.../table/metadata/00000-9441e604-b3c2-498a-a45a-6320e8ab9006.metadata.jsonnullnullnull
2022-07-28 10:43:57.487s3://.../table/metadata/00001-f30823df-b745-4a0a-b293-7532e0c99986.metadata.json17026083367764530001
2022-07-28 10:43:58.25s3://.../table/metadata/00002-2cc2837a-02dc-4687-acc1-b4d86ea486f4.metadata.json95890649397670977402

快照

显示表的有效快照:

SELECT * FROM prod.db.table$snapshots;
committed_atsnapshot_idparent_idoperationmanifest_listsummary
2019-02-08 03:29:51.21557897183625154nullappends3://.../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_atoperationsnapshot_idis_current_ancestorsummary[flink.job-id]
2019-02-08 03:29:51.215append57897183625154true2e274eecb503d85369fb390e8956c813

文件

显示表的当前数据文件:

SELECT * FROM prod.db.table$files;
contentfile_pathfile_formatspec_idpartitionrecord_countfile_size_in_bytescolumn_sizesvalue_countsnull_value_countsnan_value_countslower_boundsupper_boundskey_metadatasplit_offsetsequality_idssort_order_id
0s3:/.../table/data/00000-3-8d6d60e8-d427-4809-bcf0-f5d45a4aad96.parquetPARQUET0{1999-01-01, 01}1597[1 -> 90, 2 -> 62][1 -> 1, 2 -> 1][1 -> 0, 2 -> 0][][1 -> , 2 -> c][1 -> , 2 -> c]null[4]nullnull
0s3:/.../table/data/00001-4-8d6d60e8-d427-4809-bcf0-f5d45a4aad96.parquetPARQUET0{1999-01-01, 02}1597[1 -> 90, 2 -> 62][1 -> 1, 2 -> 1][1 -> 0, 2 -> 0][][1 -> , 2 -> b][1 -> , 2 -> b]null[4]nullnull
0s3:/.../table/data/00002-5-8d6d60e8-d427-4809-bcf0-f5d45a4aad96.parquetPARQUET0{1999-01-01, 03}1597[1 -> 90, 2 -> 62][1 -> 1, 2 -> 1][1 -> 0, 2 -> 0][][1 -> , 2 -> a][1 -> , 2 -> a]null[4]nullnull

清单(Manifests)

显示表当前的文件清单:

SELECT * FROM prod.db.table$manifests;
pathlengthpartition_spec_idadded_snapshot_idadded_data_files_countexisting_data_files_countdeleted_data_files_countpartition_summaries
s3://.../table/metadata/45b5290b-ee61-4788-b324-b1e2735c0e10-m0.avro447906668963634911763636800[[false,null,2019-05-13,2019-05-15]]

注意:

  1. manifests 表中 partition_summaries 列内的各字段对应于清单列表中的 field_summary 结构体,顺序如下:

    • contains_null
    • contains_nan
    • lower_bound
    • upper_bound
  2. contains_nan 可能返回 null,表示文件的元数据中没有该信息。这通常发生在读取 V1 表时,因为 V1 表不会填充 contains_nan。

分区

要显示表的当前分区:

SELECT * FROM prod.db.table$partitions;
partitionspec_idrecord_countfile_counttotal_data_file_size_in_bytesposition_delete_record_countposition_delete_file_countequality_delete_record_countequality_delete_file_countlast_updated_at(μs)last_updated_snapshot_id
{20211001, 11}011100210016330860341920009205185327307503337
{20211002, 11}04350011001633172537358000867027598972211003
{20211001, 10}074700000016330825987160003280122546965981531
{20211002, 10}032400001116331691594890006941468797545315876

注意:对于未分区的表,partitions 表中将不包含 partition 和 spec_id 字段。

所有元数据表

这些表是当前快照对应的各元数据表的并集,返回所有快照范围内的元数据。

危险

"all" 元数据表可能会为同一个数据文件或清单文件返回多行,因为元数据文件可能同时属于多个表快照。

所有数据文件

展示表中所有的数据文件以及每个文件的元数据:

SELECT * FROM prod.db.table$all_data_files;

所有 Manifest 文件

要显示表的所有 manifest 文件:

SELECT * FROM prod.db.table$all_manifests;
pathlengthpartition_spec_idadded_snapshot_idadded_data_files_countexisting_data_files_countdeleted_data_files_countpartition_summaries
s3://.../metadata/a85f78c5-3222-4b37-b7e4-faf944425d48-m0.avro637606272782676904868561200[{false, false, 20210101, 20210101}]

注意:

  1. manifests 表中 partition_summaries 列内的字段对应于清单列表中的 field_summary 结构体,顺序如下:

    • contains_null
    • contains_nan
    • lower_bound
    • upper_bound
  2. contains_nan 可能返回 null,表示文件的元数据中没有该信息。这种情况通常出现在读取 V1 表时,因为 V1 表不会填充 contains_nan。

参考

显示表已知的快照引用(snapshot reference):

SELECT * FROM prod.db.table$refs;
nametypesnapshot_idmax_reference_age_in_msmin_snapshots_to_keepmax_snapshot_age_in_ms
mainBRANCH4686954189838128572102030
testTagTAG468695418983812857210nullnull

评论

登录后参与评论

正在加载评论…