改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

探索Iceberg与Flink DataStream的完美集成之旅

ckckck2025年1月10日80 浏览

1 环境准备

1.1 配置pom文件

新建Maven工程,pom文件配置如下:

xml 复制代码
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>
    <groupId>org.iceberg.demo</groupId>
    <artifactId>flink-iceberg-demo</artifactId>
    <version>1.0-SNAPSHOT</version>
    <properties>
        <maven.compiler.source>8</maven.compiler.source>
        <maven.compiler.target>8</maven.compiler.target>
        <flink.version>1.16.0</flink.version>
        <java.version>1.8</java.version>
        <scala.binary.version>2.12</scala.binary.version>
        <slf4j.version>1.7.30</slf4j.version>
    </properties>
    <dependencies>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-java</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-streaming-java</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-clients</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-table-planner_${scala.binary.version}</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-connector-files</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-runtime-web</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.slf4j</groupId>
            <artifactId>slf4j-api</artifactId>
            <version>${slf4j.version}</version>
        </dependency>
        <dependency>
            <groupId>org.slf4j</groupId>
            <artifactId>slf4j-log4j12</artifactId>
            <version>${slf4j.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.logging.log4j</groupId>
            <artifactId>log4j-to-slf4j</artifactId>
            <version>2.14.0</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-statebackend-rocksdb</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.hadoop</groupId>
            <artifactId>hadoop-client</artifactId>
            <version>3.1.3</version>
        </dependency>
        <dependency>
            <groupId>org.apache.iceberg</groupId>
            <artifactId>iceberg-flink-runtime-1.16</artifactId>
            <version>1.1.0</version>
        </dependency>
    </dependencies>
</project>

1.2 配置log4j

resources目录下新建log4j.properties文件,内容如下:

properties 复制代码
log4j.rootLogger=error,stdout
log4j.appender.stdout=org.apache.log4j.ConsoleAppender
log4j.appender.stdout.target=System.out
log4j.appender.stdout.layout=org.apache.log4j.PatternLayout
log4j.appender.stdout.layout.ConversionPattern=%d %p [%c] - %m%n

2 读取数据

2.1 常规Source写法

1)Batch方式

java 复制代码
import org.apache.flink.api.common.typeinfo.Types;
import org.apache.flink.api.java.tuple.Tuple2;
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.source.FlinkSource;

public class ReadDemo {
    public static void main(String[] args) throws Exception {
        // 1. 创建Flink环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://192.168.110.120:8020/warehouse/iceberg-hadoop/iceberg_db/sample");
        DataStream<RowData> batch = FlinkSource.forRowData()
                .env(env)
                .tableLoader(tableLoader)
                .streaming(false) // false: batch方式读取数据;true: streaming方式读取数据
                .build();
        batch.map(r -> Tuple2.of(r.getInt(0), r.getString(1).toString()))
                .returns(Types.TUPLE(Types.INT, Types.STRING))
                .print();
        env.execute("Test Iceberg Read");
    }
}

2)Streaming方式

java 复制代码
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://hadoop1:8020/warehouse/spark-iceberg/default/a");
DataStream<RowData> stream = FlinkSource.forRowData()
        .env(env)
        .tableLoader(tableLoader)
        .streaming(true)
        .startSnapshotId(3821550127947089987L)
        .build();

stream.map(r -> Tuple2.of(r.getLong(0), r.getLong(1)))
        .returns(Types.TUPLE(Types.LONG, Types.LONG))
        .print();

env.execute("Test Iceberg Read");

2.2 FLIP-27 Source写法

1)Batch方式

java 复制代码
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://hadoop1:8020/warehouse/spark-iceberg/default/a");
IcebergSource<RowData> source1 = IcebergSource.forRowData()
        .tableLoader(tableLoader)
        .assignerFactory(new SimpleSplitAssignerFactory())
        .build();
DataStream<RowData> batch = env.fromSource(
        source1,
        WatermarkStrategy.noWatermarks(),
        "My Iceberg Source",
        TypeInformation.of(RowData.class));
batch.map(r -> Tuple2.of(r.getLong(0), r.getLong(1)))
        .returns(Types.TUPLE(Types.LONG, Types.LONG))
        .print();
env.execute("Test Iceberg Read");

2)Streaming方式

java 复制代码
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://hadoop1:8020/warehouse/spark-iceberg/default/a");
IcebergSource source2 = 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(
        source2,
        WatermarkStrategy.noWatermarks(),
        "My Iceberg Source",
        TypeInformation.of(RowData.class));
stream.map(r -> Tuple2.of(r.getLong(0), r.getLong(1)))
        .returns(Types.TUPLE(Types.LONG, Types.LONG))
        .print();
env.execute("Test Iceberg Read");

3 写入数据

目前支持DataStream和DataStream格式的数据流写入Iceberg表。

1)写入方式支持 append、overwrite、upsert

java 复制代码
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);

SingleOutputStreamOperator<RowData> input = env.fromElements("")
        .map(new MapFunction<String, RowData>() {
            @Override
            public RowData map(String s) throws Exception {
                GenericRowData genericRowData = new GenericRowData(2);
                genericRowData.setField(0, 99L);
                genericRowData.setField(1, 99L);
                return genericRowData;
            }
        });
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://hadoop1:8020/warehouse/spark-iceberg/default/a");

FlinkSink.forRowData(input)
        .tableLoader(tableLoader)
        .append() // append方式
        // .overwrite(true) // overwrite方式
        // .upsert(true) // upsert方式
        ;

env.execute("Test Iceberg DataStream");

写数据报错

java 复制代码
org.apache.hadoop.security.AccessControlException: Permission denied: user=Administrator, access=WRITE, inode=...

hdfs-site.xml中添加下方配置,修改后重启Hadoop:

xml 复制代码
<property>
    <name>dfs.permissions</name>
    <value>false</value>
</property>

2)写入选项

java 复制代码
FlinkSink.forRowData(input)
        .tableLoader(tableLoader)
        .set("write-format", "orc")
        .set(FlinkWriteOptions.OVERWRITE_MODE, "true");

可配置选项如下:

选项 默认值 说明
write-format Parquet 写入操作使用的文件格式:Parquet, avro或orc
target-file-size-bytes 536870912(512MB) 控制生成的文件的大小,目标大约为这么多字节
upsert-enabled false 是否启用upsert模式
overwrite-enabled false 覆盖表的数据,不能和UPSERT模式同时开启
distribution-mode none 定义写数据的分布方式:none(不打乱行)、hash(按分区键散列分布)、range(如果表有SortOrder,则通过分区键或排序键分配)
compression-codec 同 write.(fileformat).compression-codec 压缩编解码器
compression-level 同 write.(fileformat).compression-level 压缩级别
compression-strategy 同 write.orc.compression-strategy 压缩策略
java 复制代码
public class FlinkIcebergDemo1 {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // 1. 必须设置checkpoint,Flink向Iceberg中写入数据时当checkpoint发生后,才会commit数据。
        env.enableCheckpointing(5000);
        // 2. 读取Kafka中的topic数据
        KafkaSource<String> source = KafkaSource.<String>builder()
                .setBootstrapServers("192.168.6.102:6667")
                .setTopics("json")
                .setGroupId("my-group-id")
                .setStartingOffsets(OffsetsInitializer.latest())
                .setValueOnlyDeserializer(new SimpleStringSchema())
                .build();
        DataStreamSource<String> kafkaSource = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source");
        // 3. 对数据进行处理,包装成RowData对象,方便保存到Iceberg表中。
        SingleOutputStreamOperator<RowData> dataStream = kafkaSource.map(new MapFunction<String, RowData>() {
            @Override
            public RowData map(String s) throws Exception {
                System.out.println("s = " + s);
                String[] split = s.split(",");
                GenericRowData row = new GenericRowData(4);
                row.setField(0, Integer.valueOf(split[0]));
                row.setField(1, StringData.fromString(split[1]));
                row.setField(2, Integer.valueOf(split[2]));
                row.setField(3, StringData.fromString(split[3]));
                return row;
            }
        });
        // 4. 创建Hadoop配置、Catalog配置和表的Schema,方便后续向路径写数据时可以找到对应的表
        Configuration hadoopConf = new Configuration();
        Catalog catalog = new HadoopCatalog(hadoopConf, "hdfs://leidi01:8020/flinkiceberg/");
        // 配置iceberg库名和表名
        TableIdentifier name = TableIdentifier.of("icebergdb", "flink_iceberg_tbl");
        // 创建Iceberg表Schema
        Schema schema = new Schema(
                Types.NestedField.required(1, "id", Types.IntegerType.get()),
                Types.NestedField.required(2, "name", Types.StringType.get()),
                Types.NestedField.required(3, "age", Types.IntegerType.get()),
                Types.NestedField.required(4, "loc", Types.StringType.get()));
        // 如果有分区指定对应分区,这里“loc”列为分区列,可以指定unpartitioned方法不设置表分区
        PartitionSpec spec = PartitionSpec.builderFor(schema).identity("loc").build();
        // 指定Iceberg表数据格式化为Parquet存储
        Map<String, String> props = ImmutableMap.of(TableProperties.DEFAULT_FILE_FORMAT, FileFormat.PARQUET.name());
        Table table = null;
        // 通过catalog判断表是否存在,不存在就创建,存在就加载
        if (!catalog.tableExists(name)) {
            table = catalog.createTable(name, schema, spec, props);
        } else {
            table = catalog.loadTable(name);
        }
        TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://leidi01:8020/flinkiceberg//icebergdb/flink_iceberg_tbl", hadoopConf);
        // 5. 通过DataStream API向Iceberg中写入数据
        FlinkSink.forRowData(dataStream)
                .table(table)
                .tableLoader(tableLoader)
                .overwrite(false) // 默认为false,追加数据。如果设置为true就是覆盖数据
                .build();
        env.execute("DataStream API Write Data To Iceberg");
    }
}

注意事项

  1. 需要设置Checkpoint,Flink向Iceberg中写入Commit数据时,只有Checkpoint成功之后才会Commit数据,否则后期在Hive中查询不到数据
  2. 读取Kafka数据后需要包装成RowData或者Row对象,才能向Iceberg表中写出数据。写出数据时默认是追加数据,如果指定overwrite就是全部覆盖数据。
  3. 在向Iceberg表中写数据之前需要创建对应的Catalog、表Schema,否则写出时只指定对应的路径会报错找不到对应的Iceberg表。
  4. 不建议使用DataStream API向Iceberg中写数据,建议使用SQL API

4 合并小文件

Iceberg现在不支持在Flink SQL中检查表,需要使用Iceberg提供的Java API来读取元数据来获得表信息。可以通过提交Flink批处理作业将小文件重写为大文件:

java 复制代码
import org.apache.iceberg.flink.actions.Actions;

// 1. 获取Table对象
// 1.1 创建catalog对象
Configuration conf = new Configuration();
HadoopCatalog hadoopCatalog = new HadoopCatalog(conf, "hdfs://hadoop1:8020/warehouse/spark-iceberg");
// 1.2 通过catalog加载Table对象
Table table = hadoopCatalog.loadTable(TableIdentifier.of("default", "a"));
// 有Table对象,就可以获取元数据、进行维护表的操作
// 2. 通过Actions来操作合并
Actions.forTable(table)
        .rewriteDataFiles()
        .targetSizeInBytes(1024L)
        .execute();

得到Table对象,就可以获取元数据、进行维护表的操作。


大数据相关学习资料、大数据项目、湖仓一体、架构师必知必会、数据中台建设方法论...共有1400多份文档资料,另专为星球成员整理了一份比较详细的语雀知识库合计190万字和飞书文档资料。(内容太多,仅展示部分内容...),欢迎大家踊跃加入星球,您将获得:

  1. 提供最全的大数据知识库,不限设备,随时随地打开看的在线文档。
  2. 免费答疑解惑、交流技术。
  3. 面试指导、模拟面试。
  4. 各类PDF文档下载、星球代码下载。
  5. 提供简历模板,简历修改指导服务,星球成员免费提供简历修改指导。

另外说明加入星球后支持三天无理由退款,不满意无条件随时退。需要资料请加微信:D1435221412

end