Write to GCS in parquet format from apache_beam DoFn in Python

Viewed 335

I am trying to read the Pub/Sub message and write that to a GCS location in parquet format. I am following this code sample in Google documentation to have window and shards added to the file.

Any help with suggestions or code samples will be appriciated

1 Answers

The code of WriteOneFilePerWindow is defined here. It uses TextIO instead of ParquetIO.

You can follow the example to write your own PTransform with ParquetIO. Override the expand function to use

@Override
public PDone expand(PCollection<String> input) {
  ResourceId resource = FileBasedSink.convertToFileResourceIfPossible(filenamePrefix);
  FileIO.Write write = 
    FileIO.<GenericRecord>.write()
        .via(ParquetIO.sink(SCHEMA))
        .to(new PerWindowFiles(resource))
        .withTempDirectory(resource.getCurrentDirectory())
        .withWindowedWrites();
  if (numShards != null) {
    write = write.withNumShards(numShards);
  }
  return input.apply(write);
}
 

You can find doc about ParquetIO here

Related