【问题标题】:How to convince Flink to rename .inprogress files to part-xxx如何说服 Flink 将 .inprogress 文件重命名为 part-xxx
【发布时间】:2022-11-05 02:37:06
【问题描述】:

我们对流式工作流(使用 Flink 1.14.4)进行了单元测试,其中包含有界源、编写 Parquet 文件。因为它是有界的,所以会自动禁用检查点(根据 INFO msg Disabled Checkpointing. Checkpointing is not supported and not needed when executing jobs in BATCH mode.),这意味着将 ExecutionCheckpointingOptions.ENABLE_CHECKPOINTS_AFTER_TASKS_FINISH 设置为 true 无效。

是在单独的线程中运行具有无限源的线束并在没有更多数据写入输出时强制终止的唯一解决方案吗?好像很别扭...

【问题讨论】:

  • 你确定execution.checkpointing.checkpoints-after-tasks-finish.enabled 无关紧要吗?
  • 我认为在 BATCH 模式下执行有界源时,文件会自动完成。我认为不是这样吗?
  • 嗨,大卫 - 我将 execution.checkpointing.checkpoints-after-tasks-finish.enabled 设置为 true,但它并没有改变行为。但也许还有其他事情需要我解决。
  • 您使用的是 FileSink(而不是 StreamingFileSink)吗?
  • 就像您在阅读我的代码 :) 是的,我们还没有完成将所有接收器转换为新的 FileSink;一旦我们更新它,我们就会得到预期的结果。

标签: apache-flink flink-streaming


【解决方案1】:

对于其他人来说,解决方案是:

  1. 确保您使用的是FileSink,而不是旧的StreamingFileSink
  2. ExecutionCheckpointingOptions.ENABLE_CHECKPOINTS_AFTER_TASKS_FINISH 设置为真。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2014-07-26
    • 1970-01-01
    • 2011-03-02
    • 1970-01-01
    • 1970-01-01
    • 2013-08-27
    • 2011-12-08
    相关资源
    最近更新 更多