【问题标题】:Flink StreamingFileSink not ingesting to S3 when checkpointing is disabled禁用检查点时,Flink StreamingFileSink 不会摄取到 S3
【发布时间】:2020-01-22 04:33:24
【问题描述】:

我想使用 aws s3 作为 flink 中数据流的接收器。我正在使用 StreamingFileSink 类来创建接收器。

我的工作不需要检查点,但是当我禁用检查点时,数据不再写入 S3。

案例 1:启用检查点
启用检查点后,数据会成功摄取到上述 s3 路径。

案例 2:检查点禁用
禁用检查点时,数据不会写入 s3。
我尝试多次执行该作业,但每次都得到相同的结果。我在本地机器和 kubernetes 集群上都面临这个问题。

object FlinkTestJob {

  def main(args: Array[String]): Unit = {

    val env = StreamExecutionEnvironment.getExecutionEnvironment

    // with checkpointing enabled
    env.enableCheckpointing(100)

    // Sinks
    val streamStrings: Seq[String] =
      Seq("test1", "test2", "test3", "test4", "test5", "test6", "test7", "test8", "test9", "test10")

    val testStream = env.fromCollection(streamStrings)

    val rollingPolicy = new RollingPolicy[String, String] {

      override def shouldRollOnCheckpoint(partFileState: PartFileInfo[String]): Boolean =
        partFileState.getSize > 1

      override def shouldRollOnEvent(
          partFileState: PartFileInfo[String],
          element: String): Boolean = true

      override def shouldRollOnProcessingTime(
          partFileState: PartFileInfo[String],
          currentTime: Long): Boolean = true
    }

    val sink: StreamingFileSink[String] = StreamingFileSink
      .forRowFormat(new Path("s3a://testbucket/sink"), new SimpleStringEncoder[String]("UTF-8"))
      .withRollingPolicy(rollingPolicy)
      .build()

    testStream.addSink(sink)
    env.execute("test-job")
  }
}

当我使用“writeAsText("s3a://testbucket/sink")”而不是 StreamingFileSink 写入 s3 时,无论是否启用检查点,它都能正常工作。

Flink 版本:1.8.0
我想了解检查点和 StreamingFileSink 之间的关系。
谢谢

【问题讨论】:

    标签: amazon-s3 apache-flink flink-streaming


    【解决方案1】:

    要使 StreamingFileSink 工作,需要启用检查点。

    【讨论】:

      猜你喜欢
      • 2019-12-16
      • 1970-01-01
      • 2019-01-06
      • 2022-09-28
      • 2020-05-06
      • 2020-06-08
      • 1970-01-01
      • 1970-01-01
      • 2021-01-21
      相关资源
      最近更新 更多