【问题标题】:Flink Data Stream CSV Writer not writing data to CSV fileFlink Data Stream CSV Writer 未将数据写入 CSV 文件
【发布时间】:2019-01-03 07:16:44
【问题描述】:

我是 apache flink 的新手,正在尝试学习数据流。我正在从 csv 文件中读取包含 3 列(名称、主题和标记)的学生数据。我已对标记应用过滤器,并且仅选择标记> 40 的那些记录。 我正在尝试将此数据写入 csv 文件,但程序成功运行并且 csv 文件仍然为空。没有数据写入 csv 文件。

我尝试使用不同的语法来编写 csv 文件,但它们都不适合我。我通过eclipse在本地运行它。写入文本文件工作正常。

DataStream<String> text = env.readFile(format, params.get("input"), 
FileProcessingMode.PROCESS_CONTINUOUSLY,100);
DataStream<String> filtered = text.filter(new FilterFunction<String>(){
public boolean filter(String value) {
    String[] tokens = value.split(",");
    return Integer.parseInt(tokens[2]) >= 40;
}
});
filtered.writeAsText("testFilter",WriteMode.OVERWRITE);
DataStream<Tuple2<String, Integer>> tokenized = filtered
.map(new MapFunction<String, Tuple2<String, Integer>>(){
public Tuple2<String, Integer> map(String value) throws Exception {
    return new Tuple2("Test", Integer.valueOf(1));
}
});
tokenized.print(); 
tokenized.writeAsCsv("file:///home/Test/Desktop/output.csv", 
WriteMode.OVERWRITE, "/n", ",");
try {
env.execute();
} catch (Exception e1) {
e1.printStackTrace();
}
}
}

以下是我输入的 CSV 格式:

Name1,Subj1,30
Name1,Subj2,40
Name1,Subj3,40
Name1,Subj4,40

Tokenized.print() 打印所有正确的记录。

【问题讨论】:

    标签: csv apache-flink


    【解决方案1】:

    我做了一些实验,发现这个工作很好用:

    import org.apache.flink.api.java.tuple.Tuple2;
    import org.apache.flink.core.fs.FileSystem;
    import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
    
    public class WriteCSV {
        public static void main(String[] args) throws Exception {
            StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
    
            env.setParallelism(1);
    
            env.fromElements(new Tuple2<>("abc", 1), new Tuple2<>("def", 2))
                    .writeAsCsv("file:///tmp/test.csv", FileSystem.WriteMode.OVERWRITE, "\n", ",");
    
            env.execute();
        }
    }
    

    如果我不将并行度设置为 1,那么结果会有所不同。在这种情况下,test.csv 是一个包含四个文件的目录,每个文件由四个并行子任务之一编写。

    我不确定你的情况出了什么问题,但也许你可以从这个例子向后工作(假设它适合你)。

    【讨论】:

      【解决方案2】:

      您应该在tokenized.writeAsCsv(); 之前删除tokenized.print();

      它将消耗print();的数据。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2015-01-21
        • 2016-12-21
        相关资源
        最近更新 更多