【问题标题】:Fault tolerance in Flink file SinkFlink 文件 Sink 中的容错
【发布时间】:2020-04-28 14:57:57
【问题描述】:

我在集群模式下使用 Flink 流式传输与 Kafka 消费者连接器 (FlinkKafkaConsumer) 和文件接收器 (StreamingFileSink),策略只有一次。 文件接收器将文件写入本地磁盘。 我注意到,如果作业失败并且自动重新启动已打开,任务管理器会查找上次失败作业的剩余文件(隐藏文件)。 显然,由于可以将任务分配给不同的任务管理器,因此一遍又一遍地总结出更多的失败。 到目前为止,我发现的唯一解决方案是删除隐藏文件并重新提交作业。 如果我做对了(如果我错了,请纠正我),隐藏文件中的事件没有提交到引导服务器,所以没有数据丢失。

有没有办法强制 Flink 忽略已经写入的文件?或者也许有更好的方法来实现解决方案(可能以某种方式使用保存点)?

【问题讨论】:

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


    【解决方案1】:

    我在 Flink 邮件列表中得到了非常详细的答案。 TLDR,为了实现一次,我必须使用某种分布式 FS。

    完整答案:

    本地文件系统不是您想要实现的目标的正确选择。我认为您无法在此设置中实现真正的一次性政策。让我详细说明原因。 有趣的是它在检查点上的表现。该行为由 RollingPolicy 控制。由于您没有说您使用什么格式,我们假设您首先使用行格式。对于行格式,默认滚动策略(何时将文件从进行中更改为挂起)是如果文件达到 128MB、文件早于 60 秒或未写入 60 秒,它将被滚动。它不会在检查点上滚动。此外,StreamingFileSink 将文件系统视为可在还原后访问的持久接收器。这意味着它会在从检查点/保存点恢复时尝试附加到该文件。

    即使您在每个检查点滚动文件,您仍然可能会遇到可能有一些剩余的问题,因为 StreamingFileSink 在检查点完成后将文件从待处理移动到完成。如果在完成检查点和移动文件之间发生故障,它将无法在还原后移动它们(如果有访问权限,它会这样做)。

    最后,一个已完成的检查点将包含已成功处理的记录的偏移量,这意味着记录假定由 StreamingFileSink 提交。这可以是使用 StreamingFileSink 检查点元数据中的指针写入进行中文件的记录,在“待处理”文件中记录,该文件已完成的 StreamingFileSink 检查点元数据中的条目或记录在“完成”文件中。 1]

    因此,如您所见,当 StreamingFileSink 必须在重启后访问文件时,有多种情况。

    最后一件事,您提到了“提交到“引导服务器”。请记住,Flink 不会使用提交回 Kafka 的偏移量来保证一致性。它可以将这些偏移量写回,但仅用于监控/调试目的。 Flink 从其检查点存储/恢复处理后的偏移量。[3]

    如果有帮助,请告诉我。我尽力了;)顺便说一句,我强烈鼓励阅读链接的资源,因为他们试图以更有条理的方式描述所有这些。 我也在抄送 Kostas,他比我更了解 StreamingFileSink。所以他也许可以在某个地方纠正我。

    [1]https://ci.apache.org/projects/flink/flink-docs-release-1.10/dev/connectors/streamfile_sink.html [2]https://ci.apache.org/projects/flink/flink-docs-release-1.10/dev/connectors/kafka.html [3]https://ci.apache.org/projects/flink/flink-docs-release-1.10/dev/connectors/kafka.html#kafka-consumers-offset-committing-behaviour-configuration

    【讨论】:

      猜你喜欢
      • 2019-03-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-06-12
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多