【发布时间】:2019-01-08 00:01:54
【问题描述】:
我正在编写一个 flink 代码,其中我正在从本地系统读取文件并使用“writeUsingOutputFormat”将其写入数据库。
现在我的要求是写入 hdfs 而不是数据库。
你能帮我在flink中怎么做吗?
注意:hdfs 已在我的本地计算机上启动并运行。
【问题讨论】:
标签: hdfs apache-flink flink-streaming
我正在编写一个 flink 代码,其中我正在从本地系统读取文件并使用“writeUsingOutputFormat”将其写入数据库。
现在我的要求是写入 hdfs 而不是数据库。
你能帮我在flink中怎么做吗?
注意:hdfs 已在我的本地计算机上启动并运行。
【问题讨论】:
标签: hdfs apache-flink flink-streaming
Flink 提供了HDFS connector,可用于将数据写入Hadoop Filesystem 支持的任何文件系统。
提供的接收器是一个 Bucketing 接收器,它将数据流划分为包含滚动文件的文件夹。分桶行为以及写入,可以通过batch size和batch 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);
【讨论】:
在这一点上,较新的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 轻松添加。
【讨论】: