【问题标题】:Apache Fink & Iceberg: Not able to process hundred of RowData typesApache Fink & Iceberg:无法处理数百种 RowData 类型
【发布时间】:2022-10-25 18:15:30
【问题描述】:

我有一个 Flink 应用程序,它读取任意 AVRO 数据,将其映射到 RowData 并使用几个 FlinkSink 实例将数据写入 ICEBERG 表。任意数据是指我有 100 种类型的 AVRO 消息,所有这些消息都具有共同的属性“tableName”,但包含不同的列。我想将这些类型的消息中的每一种都写入一个单独的 Iceberg 表中。

为此,我使用侧输出:当我将数据映射到 RowData 时,我使用 ProcessFunction 将每条消息写入特定的 OutputTag。

稍后,随着数据流的处理,我循环进入不同的输出标签,使用 getSideOutput 获取记录,并为每个标签创建一个特定的 IcebergSink。就像是:


        final List<OutputTag<RowData>> tags = ... // list of all possible output tags

        final DataStream<RowData> rowdata = stream
                .map(new ToRowDataMap()) // Map Custom Avro Pojo into RowData
                .uid("map-row-data")
                .name("Map to RowData")
                .process(new ProcessRecordFunction(tags)) // process elements one by one sending them to a specific OutputTag
                .uid("id-process-record")
                .name("Process Input records");;

        CatalogLoader catalogLoader = ...
        String upsertField = ...
     
        outputTags
                .stream()
                .forEach(tag -> {
                    SingleOutputStreamOperator<RowData> outputStream = stream
                            .getSideOutput(tag);

                    TableIdentifier identifier = TableIdentifier.of("myDBName", tag.getId());

                    FlinkSink.Builder builder = FlinkSink
                            .forRowData(outputStream)
                            .table(catalog.loadTable(identifier))
                            .tableLoader(TableLoader.fromCatalog(catalogLoader, identifier))
                            .set("upsert-enabled", "true")
                            .uidPrefix("commiter-sink-" + tableName)
                            .equalityFieldColumns(Collections.singletonList(upsertField));
                    builder.append();
                });

当我处理几张桌子时,它工作得很好。但是当表的数量增加时,Flink 无法获取足够的任务资源,因为每个 Sink 需要两个不同的算子(因为 https://iceberg.apache.org/javadoc/0.10.0/org/apache/iceberg/flink/sink/FlinkSink.html 的内部结构)。

还有其他更有效的方法吗?或者任何优化它的方法?

提前致谢 ! :)

【问题讨论】:

    标签: apache-flink flink-streaming iceberg apache-iceberg


    【解决方案1】:

    鉴于您的问题,我假设您的操作员中约有一半是 IcebergStreamWriter ,它们已被充分利用,另一半是 IcebergFilesCommitter ,很少使用。

    您可以通过以下方式优化服务器的资源使用:

    • 增加 TaskManager 上的插槽数 (taskmanager.numberOfTaskSlots) [1] - 因此空闲 IcebergFilesCommitter Operator 未使用的 CPU 然后由 TaskManager 上的其他 Operator 使用
    • 增加提供给 TaskManager 的资源 (taskmanager.memory.process.size) [2] - 这有助于在此 TaskManager 上正在运行的 Operator 之间分配 JVM 内存开销(不要忘记并行增加插槽以开始使用额外资源:))

    为 TaskManager 添加更多插槽可能会导致 Operator 争夺 CPU,而内存仍为“空闲”任务保留。 [3]

    也许这个 Flink 架构也有用[4]

    我希望这有帮助, 彼得

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2022-06-12
      • 2020-04-19
      • 1970-01-01
      • 2015-09-09
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多