使用FileSink从Kafka源存数据时文件无法从inprogress转为finished
问题:Flink FileSink从Kafka写文件时,文件一直卡在inprogress状态转不了finished
现象总结
- 用Flink的FileSink组件将Kafka数据源的数据写入文件时,文件始终停留在
inprogress状态,无法转换为finished状态 - 替换为随机生成的流数据源时,文件能正常转为
finished状态 - 已尝试调整RollingPolicy参数、修改并行度、更改检查点间隔,问题仍未解决
排查方向
- Kafka offset提交与检查点不同步:FlinkKafkaConsumer默认的offset提交逻辑可能和FileSink的事务流程不匹配,导致文件无法完成收尾
- 反序列化或数据处理卡死:如果Kafka消息反序列化出错未被捕获,会导致数据流卡住,Sink无法正常处理后续流程
- 检查点未正常完成:FileSink的exactly-once语义依赖检查点提交事务,如果检查点失败、超时或触发不规律,文件会一直停留在inprogress状态
- 桶检查间隔过长:默认的桶检查间隔可能不够频繁,导致滚动策略无法及时触发
修复方案
1. 强制Kafka offset在检查点完成时提交
给FlinkKafkaConsumer添加配置,让offset提交和FileSink的事务流程对齐:
consumer.setCommitOffsetsOnCheckpoints(true);
2. 给数据流添加异常处理,避免卡死
检查socialStockSerializerDeserializer的实现,确保不会因坏消息抛出未捕获异常;也可以在数据流中添加异常处理环节,跳过错误数据:
DataStream<socialMediaStocks2> inputDataStream = env.addSource(consumer) .process(new ProcessFunction<socialMediaStocks2, socialMediaStocks2>() { @Override public void processElement(socialMediaStocks2 value, Context ctx, Collector<socialMediaStocks2> out) throws Exception { try { out.collect(value); } catch (Exception e) { System.err.println("处理数据出错,跳过该消息:" + e.getMessage()); } } });
3. 调整检查点配置,确保稳定触发完成
如果检查点存在超时或失败情况,添加超时时间和最小间隔配置:
env.getCheckpointConfig().setCheckpointTimeout(30000); // 设置30秒检查点超时 env.getCheckpointConfig().setMinPauseBetweenCheckpoints(5000); // 两次检查点间隔至少5秒
4. 缩短FileSink的桶检查间隔
让FileSink更频繁地检查桶状态,及时触发滚动策略:
FileSink<socialMediaStocks2> fileSink = FileSink // 其他原有配置 .withBucketCheckInterval(1000) // 每秒检查一次桶状态 .build();
修改后的完整代码示例
package org.SCAU.fileSink; import org.SCAU.SerializerDeserializer.socialStockSerializerDeserializer; import org.SCAU.model.socialMediaStocks2; import org.apache.flink.api.common.serialization.SimpleStringEncoder; import org.apache.flink.connector.file.sink.FileSink; import org.apache.flink.core.fs.Path; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.functions.ProcessFunction; import org.apache.flink.streaming.api.functions.sink.filesystem.OutputFileConfig; import org.apache.flink.streaming.api.functions.sink.filesystem.bucketassigners.DateTimeBucketAssigner; import org.apache.flink.streaming.api.functions.sink.filesystem.rollingpolicies.DefaultRollingPolicy; import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer; import org.apache.flink.util.Collector; import java.util.Properties; import java.util.concurrent.TimeUnit; public class exampleFromZhiHu { public static void main(String[] args) throws Exception{ StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(2); // 配置检查点,确保稳定触发 env.enableCheckpointing(10000); // 10秒一次检查点 env.getCheckpointConfig().setCheckpointTimeout(30000); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(5000); String intTopic = "test"; Properties properties = new Properties(); properties.setProperty("bootstrap.servers", "192.168.199.165:5092"); properties.setProperty("group.id", "test"); FlinkKafkaConsumer<socialMediaStocks2> consumer = new FlinkKafkaConsumer<>( intTopic, new socialStockSerializerDeserializer(), properties ); consumer.setStartFromLatest(); // 绑定offset提交到检查点完成 consumer.setCommitOffsetsOnCheckpoints(true); // 增加异常处理,避免坏消息卡住数据流 DataStream<socialMediaStocks2> inputDataStream = env.addSource(consumer) .process(new ProcessFunction<socialMediaStocks2, socialMediaStocks2>() { @Override public void processElement(socialMediaStocks2 value, Context ctx, Collector<socialMediaStocks2> out) throws Exception { try { out.collect(value); } catch (Exception e) { System.err.printf("处理消息异常,跳过该消息: %s%n", e.getMessage()); } } }); OutputFileConfig config = OutputFileConfig .builder() .withPartPrefix("prefix") .withPartSuffix(".txt") .build(); FileSink<socialMediaStocks2> fileSink = FileSink .forRowFormat( new Path("output/fileSinkTest"), new SimpleStringEncoder<socialMediaStocks2>("UTF-8") ) .withBucketAssigner(new DateTimeBucketAssigner<>()) .withRollingPolicy( DefaultRollingPolicy .builder() .withRolloverInterval(2000) .withInactivityInterval(1000) .withMaxPartSize(1024 * 1024 * 1024) .build() ) .withOutputFileConfig(config) .withBucketCheckInterval(1000) // 缩短桶检查间隔 .build(); inputDataStream.sinkTo(fileSink); env.execute("Kafka-To-FileSink-Job"); } }
内容的提问来源于stack exchange,提问作者HRH12138
相关产品推荐
相关产品推荐

