Sunidhi Tiwari
06/17/2023, 5:42 PMfinal StreamingFileSink<String> sink = StreamingFileSink
.forRowFormat(new Path(outputPath), new SimpleStringEncoder<String>("UTF-8"))
.withRollingPolicy(
DefaultRollingPolicy.builder()
.withRolloverInterval(TimeUnit.MINUTES.toMillis(15))
.withInactivityInterval(TimeUnit.MINUTES.toMillis(5))
.withMaxPartSize(1024 * 1024 * 1024)
.build())
.build();
The above code will assign records to the default one hour time buckets. I want to write all the parts in same directory, don't want hourly bucket. Is there is way to do so?Alokh P
06/18/2023, 5:42 AMStreamingFileSink<Tuple2<Integer, Integer>> sink = StreamingFileSink
.forRowFormat((new Path(outputPath), new SimpleStringEncoder<>("UTF-8"))
.withBucketAssigner(new KeyBucketAssigner())
.withRollingPolicy(OnCheckpointRollingPolicy.build())
.withOutputFileConfig(config)
.build();
https://nightlies.apache.org/flink/flink-docs-release-1.13/docs/connectors/datastream/streamfile_sink/#part-file-lifecycleSunidhi Tiwari
06/18/2023, 1:20 PMSunidhi Tiwari
06/18/2023, 4:07 PMSunidhi Tiwari
06/19/2023, 5:34 AMAlokh P
06/19/2023, 7:11 AMAlokh P
06/19/2023, 7:12 AMSunidhi Tiwari
06/19/2023, 8:29 AMAlokh P
06/19/2023, 9:11 AMSunidhi Tiwari
06/19/2023, 2:23 PM