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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 17:05:55