【发布时间】: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