【问题标题】:Handling output data from flink datastream处理来自 flink 数据流的输出数据
【发布时间】:2019-02-25 18:23:40
【问题描述】:

下面是我的流处理的伪代码。

Datastream env = StreamExecutionEnvironment.getExecutionEnvironment()    
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime)

Datastream stream = env.addSource() .map(mapping to java object) 
    .filter(filter for specific type of events) 
    .assignTimestampsAndWatermarks(new BoundedOutOfOrdernessTimestampExtractor(Time.seconds(2)){})
    .timeWindowAll(Time.seconds(10));

//collect all records.
Datastream windowedStream = stream.apply(new AllWindowFunction(...))

Datastream processedStream = windowedStream.keyBy(...).reduce(...)

String outputPath = ""

final StreamingFileSink sink = StreamingFileSink.forRowFormat(...).build();

processedStream.addSink(sink)

上面的代码流程是创建多个文件,我猜每个文件都有不同窗口的记录。例如,每个文件中的记录都有 30-40 秒之间的时间戳,而窗口时间只有 10 秒。 我预期的输出模式是将每个窗口数据写入单独的文件。 对此的任何参考或意见都会有很大帮助。

【问题讨论】:

  • 对于上述问题,我理解,由于我使用的是 streamingFileSink,它会将所有记录写入单个文件,并且由于我的机器是多核的,它会创建多个文件。但是我仍然需要使它能够将不同的窗口流输出写入不同的文件。对此的任何意见都会有所帮助。

标签: apache-flink flink-streaming


【解决方案1】:

看看BucketAssigner 接口。它应该足够灵活以满足您的需求。您只需要确保您的流事件包含足够的信息来确定您希望它们写入的路径。

【讨论】:

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