微信公众号搜"智元新知"关注
微信扫一扫可直接关注哦!

练习 : Flink sink to file

 

 

package test;

import bean.Stu;
import org.apache.flink.api.common.serialization.SimpleStringEncoder;
import org.apache.flink.core.fs.Path;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.sink.filesystem.StreamingFileSink;
import org.apache.flink.streaming.api.functions.sink.filesystem.rollingpolicies.DefaultRollingPolicy;
import org.apache.flink.streaming.api.functions.source.sourceFunction;import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import java.util.Random;
import java.util.concurrent.TimeUnit;

public class SinkToFile {
    public static void main(String[] args) {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);
        DataStreamSource<Stu> source = env.addSource(new SourceFunction<Stu>() {
            private boolean running = true;

            @Override
            public void run(SourceContext<Stu> sourceContext) throws Exception {
                while (running) {
                    for (int i = 0; i < 10; i++) {
                        ArrayList<String> subs = new ArrayList<String>(Arrays.asList("语文", "数学", "英语", "化学", "物理", "生物"));
                        List<String> names = Arrays.asList("张三", "李四", "王五", "赵六","田七");
                        int next = new Random().nextInt(15);
                        int random = new Random().nextInt(101);
                        sourceContext.collect(new Stu(names.get(next * random % 5), subs.get(next * random % 6), random));
                        Thread.sleep(1000);
                    }
                    Thread.sleep(20000);
                }
            }

            @Override
            public void cancel() {
                running = false;
            }
        });

        //fileSink
         final StreamingFileSink<String> fileSink = StreamingFileSink
                 .forRowFormat(new Path("./note"), new SimpleStringEncoder<String>("UTF-8"))
                 .withRollingPolicy(
                         DefaultRollingPolicy.builder()
                                 //每 15min 产生新日志文件
                                 .withRolloverInterval(TimeUnit.MINUTES.toMillis(15))
                                 //每断开 5min 产生新日志文件
                                 .withInactivityInterval(TimeUnit.MINUTES.toMillis(5))
                                 //每 1G 产生新日志文件
                                 .withMaxPartSize(1024 * 1024 * 1024)
                                 .build())
                 .build();
         //存入
        source.map(d->d.toString()).addSink(fileSink);

        try {
            env.execute();
        } catch (Exception e) {
            e.printstacktrace();
        }


        }

}

 

 

版权声明:本文内容由互联网用户自发贡献,该文观点与技术仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 [email protected] 举报,一经查实,本站将立刻删除。

相关推荐