3 ms·
> What happens is that each Spark partition outputs its own CSV, which is never what anyone would want! 1 file per partition is exactly what you want when the
by nchammas 10y ago
> What happens is that each Spark partition outputs its own CSV, which is never what anyone would want!
1 file per partition is exactly what you want when the output is large, so that multiple executors can share the work of writing out the output.
> The expected (undocumented) solution is to compress into one partition, but then non-primitive data structures get output incorrectly.
By "compress" do you mean coalesce? Yes, coalescing the RDD/DataFrame into one partition is the commonly accepted solution for when you want to force 1 output file. The cost of doing so is that you lose parallelism, since only 1 task will be able to write out that file.
And what do you mean by "non-primitive data structure" in this case?
- minimaxir 10y agoYes, I meant coalesce. It makes sense that output is distributed, but it caught me by surprise since I had not seen anything about that on the official site, nor in any literature. It is unusual since the other Spark libraries have smart defaults. My CSV had a columns of DenseVectors containing floats and when exported, it used the internal representation of numbers (see second tweet in reply to first tweet)
- nchammas 10y agoI think it would be dangerous to have the default be "coalesce to 1 partition before writing", but I agree this should be better documented since it takes many people by surprise. As for the DenseVectors, that looks strange and is perhaps worthy of a report on the project tracker.