Flink 1.20文件Sink示例:解决writeAsText等API废弃问题
Flink 1.20 废弃API替代方案(Word Count示例)
针对你代码中env.fromElements和writeAsText两个废弃API,以下是官方推荐的替代方案及修改后的完整代码:
废弃API替代说明
env.fromElements:替换为env.fromCollection(传入集合对象)或DataStreamSource.fromElements静态方法,两者均为Flink 1.20的推荐写法writeAsText:替换为FileSink连接器,这是Flink当前标准的文件输出组件,支持滚动策略、分区、文件命名等灵活配置
修改后的完整代码
import org.apache.flink.api.common.functions.FilterFunction; import org.apache.flink.api.common.functions.MapFunction; import org.apache.flink.api.java.tuple.Tuple2; import org.apache.flink.connector.file.sink.FileSink; import org.apache.flink.core.fs.Path; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.functions.sink.filesystem.OutputFileConfig; import org.apache.flink.streaming.api.functions.sink.filesystem.bucketassigners.DateTimeBucketAssigner; import org.apache.flink.streaming.api.functions.sink.filesystem.rollingpolicies.DefaultRollingPolicy; import org.apache.flink.streaming.api.functions.sink.filesystem.encoder.SimpleStringEncoder; import java.util.Arrays; import java.util.concurrent.TimeUnit; public class WordCountExample { public static void main(String[] args) { try { final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 替代废弃的env.fromElements DataStream<String> text = env.fromCollection(Arrays.asList("Nathan", "Noah", "Olivia", "Emily", "Nathaniel")); DataStream<String> filtered = text.filter(new FilterFunction<String>() { @Override public boolean filter(String value) { return value.startsWith("N"); } }); DataStream<Tuple2<String, Integer>> tokenized = filtered.map(new Tokenizer()); DataStream<Tuple2<String, Integer>> counts = tokenized .keyBy(value -> value.f0) .sum(1); DataStream<String> result = counts.map(new MapFunction<Tuple2<String, Integer>, String>() { @Override public String map(Tuple2<String, Integer> value) { return value.f0 + ": " + value.f1; } }); // 构建FileSink替代writeAsText FileSink<String> fileSink = FileSink .forRowFormat(new Path("/Users/T/Repos/FlinkPOC/src/main/java/com/poc/flink/output"), new SimpleStringEncoder<>("UTF-8")) // 按时间分区(可选) .withBucketAssigner(new DateTimeBucketAssigner<>("yyyy-MM-dd--HH")) // 滚动策略:15分钟/100MB/5分钟无数据则生成新文件(可选) .withRollingPolicy( DefaultRollingPolicy.builder() .withRolloverInterval(15, TimeUnit.MINUTES) .withInactivityInterval(5, TimeUnit.MINUTES) .withMaxPartSize(100 * 1024 * 1024) .build() ) // 自定义文件前缀后缀(可选) .withOutputFileConfig( OutputFileConfig.builder() .withPartPrefix("word-count-") .withPartSuffix(".txt") .build() ) .build(); result.sinkTo(fileSink); env.execute("Word Count Example with datastreams"); } catch (Exception e) { e.printStackTrace(); } } // 补全Tokenizer实现 public static class Tokenizer implements MapFunction<String, Tuple2<String, Integer>> { @Override public Tuple2<String, Integer> map(String value) { return new Tuple2<>(value, 1); } } }
关键注意事项
- 依赖检查:确保项目依赖中包含
flink-connector-files(Flink 1.20中通常已集成,若缺失需手动添加) - FileSink配置:示例中的分区、滚动策略等均为可选配置,可根据实际需求调整,比如不需要分区可移除
.withBucketAssigner部分 - 输出路径:注意输出路径需为目录而非具体文件名,FileSink会自动在目录下生成分区文件夹和数据文件
内容的提问来源于stack exchange,提问作者David
相关产品推荐
相关产品推荐

