Flink StreamingFileSink批量写入时如何让每个bucket的partFileIndex从0开始
问题原因
Flink 1.11版本的StreamingFileSink默认的part文件序号计数器是每个Sink子任务(Subtask)全局持有,整个任务运行周期内持续累加,不会因为切换新的bucket目录就重置,因此会出现跨小时目录序号持续递增的情况。
实现方案
方案1:自定义Bucket逻辑适配1.11版本(推荐,无需升级集群)
1.11版本没有官方提供的按bucket重置序号的开关,需要自行扩展相关实现:
- 自定义
Bucket实现类,为每个Bucket实例单独维护part文件计数器,替换默认的Subtask全局计数器 - 自定义
BucketFactory实现类,返回自定义的Bucket实例 - 在构造
StreamingFileSink时传入自定义的BucketFactory即可
核心代码示例如下:
import org.apache.flink.streaming.api.functions.sink.filesystem.{Bucket, BucketFactory, OutputFileConfig} import org.apache.hadoop.fs.Path // 自定义Bucket,每个Bucket维护独立的计数器 class PerBucketCounterBucket( bucketId: String, bucketPath: Path, subtaskIndex: Int, outputFileConfig: OutputFileConfig ) extends Bucket[RawSample, String, String](bucketId, bucketPath, subtaskIndex, outputFileConfig) { // 当前Bucket专属的part计数器 private var partCounter: Long = 0L override def getPartFileName(partId: Long): String = { val currentCounter = partCounter partCounter += 1 s"${outputFileConfig.getPartPrefix}-$subtaskIndex-$currentCounter${outputFileConfig.getPartSuffix}" } // 重写快照逻辑,将计数器纳入状态管理,避免故障恢复后序号重置 override def onSnapshotState(checkpointId: Long): Unit = { super.onSnapshotState(checkpointId) // 此处将partCounter写入Bucket的状态快照即可 } // 重写恢复逻辑,读取快照中存储的计数器值 override def initializeState(checkpointId: Long): Unit = { super.initializeState(checkpointId) // 此处读取快照中的partCounter赋值给当前变量即可 } } // 自定义Bucket工厂,生成自定义Bucket实例 class PerBucketCounterBucketFactory extends BucketFactory[RawSample, String, String] { override def getNewBucket( subtaskIndex: Int, bucketId: String, bucketPath: Path, creationTime: Long, outputFileConfig: OutputFileConfig ): Bucket[RawSample, String, String] = { new PerBucketCounterBucket(bucketId, bucketPath, subtaskIndex, outputFileConfig) } }
修改你的Sink构造代码,添加自定义Bucket工厂配置:
val sinker = StreamingFileSink .forBulkFormat(new Path(option.dumpOutputPath), writer) .withBucketAssigner(new DateTimeBucketAssigner[RawSample]("yyyy-MM-dd/HH")) .withRollingPolicy(OnCheckpointRollingPolicy.build()) .withBucketCheckInterval(option.rolloverInterval) .withOutputFileConfig(OutputFileConfig.builder().withPartSuffix(".gz.parquet").build()) // 新增本行配置,传入自定义Bucket工厂 .withBucketFactory(new PerBucketCounterBucketFactory()) .build()
方案2:升级Flink版本(可选)
如果可以升级集群Flink版本到1.13及以上,官方已经提供了开箱支持的配置项,只需要在Flink配置中将execution.sink.file.part-counter-reset-on-new-bucket设置为true,不需要修改业务代码即可实现每个Bucket下的part序号从0开始计数。
内容的提问来源于stack exchange,提问作者kemiya
相关产品推荐
相关产品推荐

