【问题标题】:Regarding flink stream sink to hdfs关于 flink 流 sink 到 hdfs
【发布时间】:2019-01-08 00:01:54
【问题描述】:

我正在编写一个 flink 代码,其中我正在从本地系统读取文件并使用“writeUsingOutputFormat”将其写入数据库。

现在我的要求是写入 hdfs 而不是数据库。

你能帮我在flink中怎么做吗?

注意:hdfs 已在我的本地计算机上启动并运行。

【问题讨论】:

    标签: hdfs apache-flink flink-streaming


    【解决方案1】:

    Flink 提供了HDFS connector,可用于将数据写入Hadoop Filesystem 支持的任何文件系统。

    提供的接收器是一个 Bucketing 接收器,它将数据流划分为包含滚动文件的文件夹。分桶行为以及写入,可以通过batch sizebatch roll over time interval等参数进行配置

    Flink 文档给出了以下示例 -

    DataStream<Tuple2<IntWritable,Text>> input = ...;
    
    BucketingSink<String> sink = new BucketingSink<String>("/base/path");
    sink.setBucketer(new DateTimeBucketer<String>("yyyy-MM-dd--HHmm", ZoneId.of("America/Los_Angeles")));
    sink.setWriter(new SequenceFileWriter<IntWritable, Text>());
    sink.setBatchSize(1024 * 1024 * 400); // this is 400 MB,
    sink.setBatchRolloverInterval(20 * 60 * 1000); // this is 20 mins
    
    input.addSink(sink);
    

    【讨论】:

      【解决方案2】:

      在这一点上,较新的Streaming File Sink 可能是比 Bucketing Sink 更好的选择。此描述来自 Flink 1.6 发行说明(注意 Flink 1.7 中添加了对 S3 的支持):

      新的 StreamingFileSink 是一个完全一次性的接收器,用于写入 利用从 以前的 BucketingSink。通过集成支持 Exactly-once 使用 Flink 的检查点机制的 sink。新水槽是 建立在 Flink 自己的 FileSystem 抽象之上,它支持本地 文件系统和 HDFS,计划在不久的将来支持 S3 [现在包含在 Flink 1.7 中]。它 公开可插入的文件滚动和分桶策略。除了 逐行编码格式,新的 StreamingFileSink 附带 支持实木复合地板。其他批量编码格式,如 ORC 使用公开的 API 轻松添加。

      【讨论】:

      • 我同意,对于无限制的流式传输场景,您认为 StreamingFileSink 是首选是正确的。上面写着:“从本地系统读取文件”。在有界流场景中,StreamingFileSink 依赖于检查点来最终确定输出文件(而不是在流关闭时)。这可能会导致事件不被写入下沉并丢失。这有待解决此问题:issues.apache.org/jira/browse/FLINK-2646
      猜你喜欢
      • 2019-03-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-08-20
      • 1970-01-01
      相关资源
      最近更新 更多