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