class PartitionTypeBucketAssigner( private val partitionType: PartitionColumnType, private val dateFormat: String ) : DateTimeBucketAssigner() { private val formatter = DateTimeFormatter.ofPattern(dateFormat) override fun getBucketId(element: MyProtobufModel, context: BucketAssigner.Context): String { return when (partitionType) { PartitionColumnType.EVENT_TIME -> { val eventTime = Instant.ofEpochMilli(element.timestamp.getSeconds() * 1000) formatter.format(eventTime) } PartitionColumnType.PROCESSING_TIME -> { val contextProcessingTime = context.currentProcessingTime() val processingTime = Instant.ofEpochMilli(contextProcessingTime) formatter.format(processingTime) } } } } class MyProtobufModelToParquet : FlinkApp() { val partitionDateTimeBucketAssigner = PartitionTypeBucketAssigner( partitionType = s3PartitionColumn, dateFormat = "'year='yyyy/'month='M/'day='d" ) override fun output(outputStream: DataStream) { val sink:StreamingFileSink = StreamingFileSink .forBulkFormat( Path("s3://test-bucket/test-path/"), ParquetProtoWriters.forType(MyProtobufModel::class.java)) .withRollingPolicy( OnCheckpointRollingPolicy.build()) .withBucketAssigner(partitionDateTimeBucketAssigner) .build() outputStream.addSink(sink) }