Cannot write parquet to amazon s3 bucket using AvroParquetWriter in Java

Viewed 16

Hi am trying to write parquets to an amazon s3 bucket/key using JAVA but getting an [org.apache.hadoop.fs.UnsupportedFileSystemException: No FileSystem for scheme "file"].

Cannot write [s3a://nprd-pr-snd-rtgrev-edp/sndrtgrev/out/jcl-0.snappy.parquet]: org.apache.hadoop.fs.UnsupportedFileSystemException: No FileSystem for scheme "file" at org.apache.hadoop.fs.FileSystem.getFileSystemClass(FileSystem.java:3443) at org.apache.hadoop.fs.FileSystem.createFileSystem(FileSystem.java:3466) at org.apache.hadoop.fs.FileSystem.access$300(FileSystem.java:174) at org.apache.hadoop.fs.FileSystem$Cache.getInternal(FileSystem.java:3574) at org.apache.hadoop.fs.FileSystem$Cache.get(FileSystem.java:3521) at org.apache.hadoop.fs.FileSystem.get(FileSystem.java:540) at org.apache.hadoop.fs.FileSystem.getLocal(FileSystem.java:496) at org.apache.hadoop.fs.LocalDirAllocator$AllocatorPerContext.confChanged(LocalDirAllocator.java:316) at org.apache.hadoop.fs.LocalDirAllocator$AllocatorPerContext.getLocalPathForWrite(LocalDirAllocator.java:393) at org.apache.hadoop.fs.LocalDirAllocator.getLocalPathForWrite(LocalDirAllocator.java:165) at org.apache.hadoop.fs.LocalDirAllocator.getLocalPathForWrite(LocalDirAllocator.java:146) at org.apache.hadoop.fs.s3a.S3AFileSystem.createTmpFileForWrite(S3AFileSystem.java:1019) at org.apache.hadoop.fs.s3a.S3ADataBlocks$DiskBlockFactory.create(S3ADataBlocks.java:816) at org.apache.hadoop.fs.s3a.S3ABlockOutputStream.createBlockIfNeeded(S3ABlockOutputStream.java:204) at org.apache.hadoop.fs.s3a.S3ABlockOutputStream.(S3ABlockOutputStream.java:182) at org.apache.hadoop.fs.s3a.S3AFileSystem.create(S3AFileSystem.java:1369) at org.apache.hadoop.fs.FileSystem.create(FileSystem.java:1195) at org.apache.hadoop.fs.FileSystem.create(FileSystem.java:1175) at org.apache.hadoop.fs.FileSystem.create(FileSystem.java:1064) at org.apache.parquet.hadoop.ParquetFileWriter.(ParquetFileWriter.java:244) at org.apache.parquet.hadoop.ParquetWriter.(ParquetWriter.java:273) at org.apache.parquet.hadoop.ParquetWriter$Builder.build(ParquetWriter.java:494) at com.jclconsultants.parquet.tool.AwsParquetProcessor.writeParquetWithAvroToAwsBucket(AwsParquetProcessor.java:210) at com.jclconsultants.parquet.tool.AwsParquetProcessor.run(AwsParquetProcessor.java:81) at com.jclconsultants.parquet.tool.ToolProcessParquetsFromAwsBucket.launchAwsParquetProcessor(ToolProcessParquetsFromAwsBucket.java:269) at com.jclconsultants.parquet.tool.ToolProcessParquetsFromAwsBucket.launchThreads(ToolProcessParquetsFromAwsBucket.java:251) at com.jclconsultants.parquet.tool.ToolProcessParquetsFromAwsBucket.processParquetsFromBucket(ToolProcessParquetsFromAwsBucket.java:160) at com.jclconsultants.parquet.tool.ToolProcessParquetsFromAwsBucket.main(ToolProcessParquetsFromAwsBucket.java:127)

Here is how I am writing (code was gathered from different post):

private void writeParquetWithAvroToAwsBucket(String outputParquetName, List<GenericData.Record> records) {
    URI awsURI = null;
    try {
        awsURI = new URI("s3a://" + bucketName + "/" + outputParquetName ); 
    } catch (URISyntaxException e1) {
        e1.printStackTrace();
        return;
    }
    Path dataFile = new Path(awsURI);

    Configuration config = new Configuration();
    config.set("fs.s3a.access.key", accesskey);
    config.set("fs.s3a.secret.key", secretAccessKey);
    config.set("fs.s3a.endpoint", "s3." + Regions.CA_CENTRAL_1.getName() + ".amazonaws.com");
    config.set("fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem");
    config.set("fs.s3a.aws.credentials.provider", "org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider");
    config.set("fs.s3a.server-side-encryption-algorithm", S3AEncryptionMethods.SSE_KMS.getMethod());
    config.set("fs.s3a.connection.ssl.enabled", "true");
    config.set("fs.s3a.impl.disable.cache", "true");
    config.set("fs.s3a.path.style.access", "true");

    try (ParquetWriter<GenericData.Record> writer = AvroParquetWriter.<GenericData.Record>builder(dataFile)//
            .withSchema(getSchema())//
            .withConf(config)//
            .withCompressionCodec(SNAPPY)//
            .withWriteMode(OVERWRITE)//
            .build()) {

        for (GenericData.Record record : records) {
            writer.write(record);
        }
    } catch (Exception e) {
        e.printStackTrace();
    } 
}

and the schema si as follow:

private Schema getSchema() {
        String json = "{\"type\":\"record\",\r\n" + " \"name\":\"spark_schema\",\r\n" + " \"fields\":[\r\n"
                + " {\"name\":\"PolicyVersion_uniqueId\",\"type\":[\"null\",\"string\"],\"default\":null},\r\n"
                + " {\"name\":\"agreementNumber\",\"type\":[\"null\",\"string\"],\"default\":null},\r\n"
                + " {\"name\":\"epoch_time\",\"type\":[\"null\",\"string\"],\"default\":null},\r\n"
                + " {\"name\":\"xml\",\"type\":[\"null\",\"string\"],\"default\":null},\r\n"
                + " {\"name\":\"filename\",\"type\":[\"null\",\"string\"],\"default\":null},\r\n"
                + " {\"name\":\"message_header\",\"type\":[\"null\",{\"type\":\"map\",\"values\":[\"null\",\"string\"]}],\"default\":null},\r\n"
                + " {\"name\":\"tracking_number\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"fullTermPremium\",\"type\":[\"null\",\"string\"],\"default\":null},\r\n"
                + " {\"name\":\"epoch_date\",\"type\":[\"null\",{\"type\":\"int\",\"logicalType\":\"date\"}],\"default\":null},\r\n"
                + " {\"name\":\"collection_date\",\"type\":[\"null\",{\"type\":\"int\",\"logicalType\":\"date\"}],\"default\":null}\r\n"
                + "]}";
        return new Schema.Parser().parse(json);
    }

I've also tried with a different credential provider (AnonymousAWSCredentialsProvider and TemporaryAWSCredentialsProvider) but still did not get it to work. I have a feeling that my config is missing something as it fails during the AvroParquetWriter build().

What am I doing wrong?

Notice that I can read parquets from S3 with a similar configuration.
I can also write a parquet file to my local drive and then upload the parquet file to s3 as follow:

private Path writeGenericRecordsToLocalDrive() {
    Path dataFile = new Path(OUTPUT_DIR + "/" + parquetName);
    Configuration config = new Configuration();
    config.set("fs.hdfs.impl", org.apache.hadoop.hdfs.DistributedFileSystem.class.getName());
    config.set("fs.file.impl", org.apache.hadoop.fs.LocalFileSystem.class.getName());

    try (ParquetWriter<GenericData.Record> writer = AvroParquetWriter.<GenericData.Record>builder(dataFile)//
            .withSchema(getSchema())//
            .withConf(config)//
            .withCompressionCodec(SNAPPY)//
            .withWriteMode(OVERWRITE)//
            .build()) {

        for (GenericData.Record parquet : parquetRecords) {
            writer.write(parquet);
            nbRecordsInParquet++;
        }
    } catch (IOException e) {
        e.printStackTrace();
    }
    return dataFile;
}

private void uploadLocalParquetToAwsBucket(Path dataFile) {
    AmazonS3 s3Client = null;
    try {
        File fileToUpload = new File(OUTPUT_DIR + "/" + dataFile.getName());

        ObjectMetadata objectMetadata = new ObjectMetadata();
        objectMetadata.setHeader(Headers.SERVER_SIDE_ENCRYPTION, "aws:kms");
        objectMetadata.setContentLength(fileToUpload.length());

        BasicAWSCredentials creds = new BasicAWSCredentials(accessKey, secretAccessKey);
        s3Client = AmazonS3ClientBuilder.standard().withRegion(Regions.CA_CENTRAL_1)
                .withCredentials(new AWSStaticCredentialsProvider(creds)).build();

        String bucketDestination = fileToUpload.getName();
        com.amazonaws.services.s3.model.PutObjectRequest putRequest = new com.amazonaws.services.s3.model.PutObjectRequest(
                bucketName + "/" + bucketFolder, bucketDestination, new FileInputStream(fileToUpload),
                objectMetadata);
        putRequest.putCustomRequestHeader(Headers.SERVER_SIDE_ENCRYPTION, "aws:kms");
        PutObjectResult putResult = s3Client.putObject(putRequest);
    } catch (Exception e) {
        e.printStackTrace();
    } finally {
        if (s3Client != null) {
            s3Client.shutdown();
        }
    }
}
0 Answers
Related