You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink Table API作业中.inprogress文件无法滚动生成为最终Parquet文件问题求助

大家好,我最近在做一个Flink流处理作业,用Table API从Kafka读取日志数据,然后写入本地路径(用来模拟S3)的Parquet文件。现在遇到一个头疼的问题:作业确实会在目标目录生成.inprogress后缀的临时文件,但这些文件始终无法最终化为正式的.parquet文件,试了好几种方法都没解决,想请各位帮忙分析下原因。

我的简化代码如下:

// 创建Kafka源表
tableEnv.executeSql( 
    "CREATE TABLE logs_source (" + 
    " `timestamp` STRING," + 
    " user_id INT," + 
    " message STRING" + 
    ") WITH (" + 
    " 'connector' = 'kafka'," + 
    " 'topic' = 'kafka-5s'," + 
    " 'properties.bootstrap.servers' = 'localhost:9092'," + 
    " 'properties.group.id' = 'kafka-5s-1'," + 
    " 'scan.startup.mode' = 'earliest-offset'," + 
    " 'format' = 'json'," + 
    " 'json.timestamp-format.standard' = 'ISO-8601'" + 
    ")" 
);

// 创建FileSystem Sink表(本地路径模拟S3)
tableEnv.executeSql( 
    "CREATE TABLE s3_sink (" + 
    " `timestamp` STRING," + 
    " user_id INT," + 
    " message STRING" + 
    ") WITH (" + 
    " 'connector' = 'filesystem'," + 
    " 'path' = 'file:///D://flink-data/data'," + 
    " 'format' = 'parquet'," + 
    " 'parquet.compression' = 'SNAPPY'," + 
    " 'sink.rolling-policy.rollover-interval' = '1000'," + 
    " 'sink.rolling-policy.file-size' = '10'," + 
    " 'sink.rolling-policy.inactivity-interval' = '1000'" + 
    ")" 
);

// 插入数据
tableEnv.executeSql("INSERT INTO s3_sink SELECT `timestamp`, user_id, message FROM logs_source");

问题现象:

作业运行后,目标目录会出现类似part-0-0.inprogress的临时文件,但这些文件永远不会完成滚动操作,变成带.parquet后缀的最终文件。

我已经尝试过的操作:

  • 把滚动策略的参数设置得非常激进:rollover-interval和inactivity-interval设为1000ms,file-size设为10字节,本以为能快速触发文件滚动,但完全没效果
  • 开启了Checkpointing,配置了文件系统作为Checkpoint存储,具体代码如下:
Configuration config = GlobalConfiguration.loadConfiguration( 
    "D:\\Gitlab\\java-project\\src\\main\\resources\\" 
);
config.set(CheckpointingOptions.CHECKPOINT_STORAGE, "filesystem");
config.set(CheckpointingOptions.CHECKPOINTS_DIRECTORY, "file:///D://flink-data/checkpoints/");

// 打印S3配置(测试阶段暂未用真实S3)
System.out.println("fs.s3a.access.key: " + config.getString("fs.s3a.access.key", "NOT SET"));
System.out.println("fs.s3a.secret.key: " + config.getString("fs.s3a.secret.key", "NOT SET"));
System.out.println("fs.s3a.endpoint: " + config.getString("fs.s3a.endpoint", "NOT SET"));

FileSystem.initialize(config, null);

// 初始化流执行环境并开启Checkpoint
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(config);
env.enableCheckpointing(10000);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(5000);
env.getCheckpointConfig().setCheckpointTimeout(60000);
env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);
env.configure(config);

// 创建TableEnvironment
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
  • 确认Kafka的kafka-5s主题一直在持续产生数据,不存在断流的情况
  • 目前用本地file:///路径做测试,还没切换到真实S3,排除了S3权限、端点配置等问题

想请教各位,是不是我哪里配置错了?比如Checkpoint的设置有没有问题?或者FileSystem Connector的滚动策略还有其他需要注意的参数?麻烦大家帮忙看看,谢谢!

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.08 10:24:38