Apache Flink Table API仅生成part文件,如何输出标准CSV文件?
问题
我尝试使用Apache Flink的Table API将流数据保存为CSV文件,但仅生成带.part前缀的临时文件,而非标准CSV文件。使用的SQL代码如下:
CREATE TABLE MyTable3 ( transactionId STRING, timestamps STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'financial_transactions', 'properties.bootstrap.servers' = 'kafka:29092', 'properties.group.id' = 'testGroup', 'scan.startup.mode' = 'earliest-offset' 'format' = 'json' ); CREATE TABLE fs_table2 ( transactionId STRING ) WITH ( 'connector'='filesystem', 'path'='abfs://project@r.dfs.core.windows.net/data-project/data/output', 'format'='csv', 'sink.rolling-policy.file-size'='128MB', 'sink.rolling-policy.rollover-interval'='2 min' ); INSERT INTO fs_table2 SELECT transactionId FROM MyTable3;
生成的文件名为类似.part-0c9f9243-f9cc-4671-b734-0ba4185adf6e-0-1.inprogress.14ad1c71-15ec-470b-ac83-98a781e3afe4的格式,请问如何生成真正的CSV文件?
解决方案
1. 先修正代码语法错误
你的MyTable3定义中,'scan.startup.mode' = 'earliest-offset'行末尾缺少逗号,会导致语法解析失败,先补上:
CREATE TABLE MyTable3 ( transactionId STRING, timestamps STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'financial_transactions', 'properties.bootstrap.servers' = 'kafka:29092', 'properties.group.id' = 'testGroup', 'scan.startup.mode' = 'earliest-offset', 'format' = 'json' );
2. 理解Flink流式文件Sink的文件生成逻辑
你看到的.inprogress后缀文件是正在写入的临时文件,这是Flink流处理的正常行为:
- 当文件满足你配置的滚动条件(达到128MB大小,或超过2分钟),Flink会自动将其重命名为不带
.inprogress的最终文件,此时文件内容就是标准CSV格式,只是文件名仍以part开头。 - 如果作业一直在运行,未触发滚动条件的文件会一直处于
.inprogress状态,手动停止作业后,所有临时文件都会被最终化为完整的CSV文件。
3. 自定义文件名(可选)
如果希望文件名更贴近“标准CSV文件”格式,可以添加两个配置项:
sink.file-name-prefix:自定义文件前缀sink.file-name-suffix:添加.csv后缀
修改后的fs_table2定义如下:
CREATE TABLE fs_table2 ( transactionId STRING ) WITH ( 'connector'='filesystem', 'path'='abfs://project@r.dfs.core.windows.net/data-project/data/output', 'format'='csv', 'sink.rolling-policy.file-size'='128MB', 'sink.rolling-policy.rollover-interval'='2 min', 'sink.file-name-prefix'='transaction_', 'sink.file-name-suffix'='.csv' );
配置后生成的文件会是类似transaction_xxxx.csv的格式,内容完全符合CSV规范。
4. 关于单个CSV文件的说明
流处理场景下,Flink无法生成单个无限大的CSV文件,因为流数据是持续输入的,只能通过滚动策略生成多个分片文件。如果必须得到单个文件,需要切换到批处理模式处理数据。
内容的提问来源于stack exchange,提问作者Raz
相关产品推荐
相关产品推荐

