You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink 1.20文件Sink示例:解决writeAsText等API废弃问题

针对你代码中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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 09:20:22