Can we decouple s3 sink commit time from checkpoint interval time?

Viewed 61

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();
0 Answers
Related