【问题标题】:Flink 1.6 bucketing sink HDFS files stuck in .in-progressFlink 1.6 bucketing sink HDFS 文件卡在 .in-progress 中
【发布时间】:2019-03-24 17:26:55
【问题描述】:

我正在将 Kafka 数据流写入 HDFS 路径中的分桶接收器。 Kafka 给出字符串数据。使用 FlinkKafkaConsumer010 从 Kafka 消费

-rw-r--r--   3 ubuntu supergroup    4097694 2018-10-19 19:16 /streaming/2018-10-19--19/_part-0-1.in-progress
-rw-r--r--   3 ubuntu supergroup    3890083 2018-10-19 19:16 /streaming/2018-10-19--19/_part-1-1.in-progress
-rw-r--r--   3 ubuntu supergroup    3910767 2018-10-19 19:16 /streaming/2018-10-19--19/_part-2-1.in-progress
-rw-r--r--   3 ubuntu supergroup    4053052 2018-10-19 19:16 /streaming/2018-10-19--19/_part-3-1.in-progress

只有当我使用一些映射函数来动态操作流数据时才会发生这种情况。如果我直接将流写入 HDFS,它工作正常。知道为什么会发生这种情况吗?我正在使用 Flink 1.6.1、Hadoop 3.1.1 和 Oracle JDK1.8

【问题讨论】:

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


    【解决方案1】:

    我遇到了类似的问题,启用检查点将状态后端从默认的 MemoryStateBackend 更改为 FsStateBackend 有效。就我而言,检查点失败是因为 MemoryStateBackendmaxStateSize 太小,以至于其中一个操作的状态无法放入内存。

    StateBackend stateBackend = new FsStateBackend("file:///home/ubuntu/flink_state_backend");
    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment()
        .enableCheckpointing(Duration.ofSeconds(60).toMillis())
        .setStateBackend(stateBackend);
    

    【讨论】:

      【解决方案2】:

      这个问题有点晚了,但我也遇到了类似的问题。 我有一个案例类地址

      case class Address(val i: Int)
      

      我从集合中读取源代码,例如地址数

          env.fromCollection(Seq(new Address(...), ...)) 
      
          ...
          val customAvroFileSink = StreamingFileSink
            .forBulkFormat(
              new Path("/tmp/data/"),
              ParquetAvroWriters.forReflectRecord(classOf[Address]))
            .build()
          ... 
          xxx.addSink(customAvroFileSink)
      

      启用检查点后,我的镶木地板文件也将在进行中结束

      我发现 Flink 在触发检查点之前完成了该过程,因此我的结果从未完全刷新到磁盘。在我将检查点间隔更改为较小的数字后,镶木地板不再进行中。

      【讨论】:

        【解决方案3】:

        这种情况通常在检查点被禁用时发生。

        您能否在使用映射功能运行作业时检查检查点设置?看起来您已为直接写入 HDFS 的作业启用检查点。

        【讨论】:

        • 我将 env 配置为 env.enableCheckpointing(5 * 1000 * 60);,但我得到的输出是 .part-0-138.inprogress.91eaf4ce-5385-46cd-b6ac-3d1b27c5b550
        猜你喜欢
        • 1970-01-01
        • 2020-05-23
        • 1970-01-01
        • 1970-01-01
        • 2014-11-27
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2019-12-10
        相关资源
        最近更新 更多