I am using Flink BulkFormatBuilder with ParquetWriter factory and I want to partition my s3 sink output at minute level. I have created custom DateTimeBucketAssigner with file format "yyyy-MM-dd-HH-mm".
Flink currently committing s3 file on every checkpoint interval and we set this interval to 1 min so that output gets available after every 1 min in s3 buckets
We had seen in past that if there was backpressure in job, this checkpoint interval got shifted by few mins, which means there are chances that our output might get delayed by few mins. We wanted to ensure our data is getting pushed on s3 bucket after every 1 min. Can we decouple this S3 commit time from checkpoint interval time?
FileSink.forBulkFormat(s3SinkPath, new ParquetWriterFactory<>(builder))
.withRollingPolicy(new CustomCheckpointFileRollingPolicy(60000))
.withBucketCheckInterval(bucketCheckInterval)
.withBucketAssigner(new CustomDateTimeBucketAssigner())
.build();