I'm attempting to read in a large csv file, do some processing and write out the data to a parquet file.
I am able to read the CSV contents, but when sending the data to a parquet sink it does not write, the akka graph finished before anything can be written.
I've narrowed down that when adding data to the GenericRecord Object (to write to parquet) nothing gets logged out. I've also tried writing to the parquet file using a flow (writer.write(record)) and ignoring the sink (with Sink.ignore()) and am still facing the same issue. I'm not sure where I am going wrong, any guidance would be appreciated.
Below is a sample of my code:
String parquetOut = "out.parquet";
String schemaFilePath = "schema.json";
final String path = "data.csv";
final java.nio.file.Path file = Paths.get(path);
ActorSystem<Void> actorSystem = ActorSystem.create(Behaviors.empty(), "actorSystem");
// ********** Start: Avro Parquet *****************
Configuration conf = new Configuration();
conf.setBoolean(AvroReadSupport.AVRO_COMPATIBILITY, true);
Schema schema = getSchema(schemaFilePath);
ParquetWriter<GenericRecord> writer =
AvroParquetWriter.<GenericRecord>builder(new Path(parquetOut))
.withConf(conf)
.withWriteMode(ParquetFileWriter.Mode.OVERWRITE)
.withCompressionCodec(CompressionCodecName.SNAPPY)
.withDictionaryEncoding(true)
.withSchema(schema)
.build();
// ********** End: Avro Parquet *****************
Source<ByteString, CompletionStage<IOResult>> source = FileIO.fromPath(file);
Flow<ByteString, GenericRecord, NotUsed> parquetFlow = Flow.of(ByteString.class).map(value -> {
List<String> values = Arrays.asList(value.utf8String().split(","));
GenericRecord record = new GenericData.Record(schema);
record.put("first-name", values.get(0));
record.put("last-name", values.get(1));
record.put("id", Integer.valueOf(values.get(2)));
record.put("age", Integer.valueOf(values.get(3)));
return record;
});
Sink<GenericRecord, CompletionStage<Done>> sink = AvroParquetSink.create(writer);
CompletionStage<IOResult> result = source.via(parquetFlow).to(sink).run(actorSystem);
result.whenComplete((input, exception) -> {
if (exception != null) {
System.out.println("exception occurred");
System.err.println(exception);
} else {
try {
writer.close();
} catch (IOException ioException) {
ioException.printStackTrace();
}
System.out.println("no exception, got result: " + input);
actorSystem.terminate();
}
});