【问题标题】:flink streaming create file ( csv or text) on time windowflink 流在时间窗口上创建文件( csv 或文本)
【发布时间】:2017-10-26 01:11:57
【问题描述】:

我是 flink 新手 我有这样的转变假设

val supportTask= customSource

  .map( line => line.split(","))
  .map( line => SupportTaskNew(line(0)toInt,line(1).toString,line(2)toString,line(3)toLong,line(4).toString,line(5)toInt,line(6)toInt))
  .filter(_ => true) //todo put sent date condition
  .map( line => Count(1))
  .keyBy(0)
  .timeWindow(Time.seconds(20)) //todo for time being 10 seconds, actuals 30 min
 .sum(0)  

现在我想为每 20 秒的时间窗口创建一个文件

supportTask.writeAsText(("D://myfile_"+Calendar.getInstance().get(Calendar.SECOND)),WriteMode.NO_OVERWRITE).setParallelism(1)

我提供了文件名+秒数,这样每次创建文件时都会附加秒数。

但是这里只创建了一个文件,我想每 20 秒创建一个新文件,我该怎么做?

【问题讨论】:

  • 使用 DataStream.writeUsingOutputFormat() APIwriteAsText 将所有输出记录写入参数中指定的文件。你必须实现一种特殊的输出格式来实现这一点。

标签: apache-kafka apache-flink flink-streaming flink-cep


【解决方案1】:

也许您可以使用Bucketing File Sink 和自定义Bucketer 来做到这一点。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-06-06
    • 2019-01-31
    • 2017-10-18
    • 2018-06-02
    • 1970-01-01
    相关资源
    最近更新 更多