Flink Table API作业中.inprogress文件无法滚动生成为最终Parquet文件问题求助
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
相关产品推荐
相关产品推荐

