news 2026/8/3 19:29:16

CSV Format Flink / PyFlink 读写 CSV 的正确姿势(含 Schema 高级配置)

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
CSV Format Flink / PyFlink 读写 CSV 的正确姿势(含 Schema 高级配置)

1、依赖引入

Java/Scala 工程需要加 Flink CSV 依赖:

<dependency><groupId>org.apache.flink</groupId><artifactId>flink-csv</artifactId><version>2.2.0</version></dependency>

PyFlink 用户一般可以直接在作业里使用(前提是集群环境里对应的 jar 能被加载;如果你是在远程集群跑,仍然需要按你前面“依赖管理”章节的方式把 jar 加入pipeline.jarsenv.add_jars())。

2、Java:快速读取 POJO(自动推导 Schema)

最省事的方式:让 Jackson 根据 POJO 字段推导 CSV schema:

CsvReaderFormat<SomePojo>csvFormat=CsvReaderFormat.forPojo(SomePojo.class);FileSource<SomePojo>source=FileSource.forRecordStreamFormat(csvFormat,Path.fromLocalFile(...)).build();

注意:CSV 列顺序必须和 POJO 字段顺序一致。必要时加:

@JsonPropertyOrder({"field1","field2",...})

否则列对不上会出现解析错位(最常见的“字段都不为空但值都错”的隐性 bug)。

3、Java:高级配置(自定义分隔符、禁用引号等)

需要精细控制时,用forSchema(...)自己生成CsvSchema,例如把分隔符改成|,并禁用 quote:

Function<CsvMapper,CsvSchema>schemaGenerator=mapper->mapper.schemaFor(CityPojo.class).withoutQuoteChar().withColumnSeparator('|');CsvReaderFormat<CityPojo>csvFormat=CsvReaderFormat.forSchema(()->newCsvMapper(),schemaGenerator,TypeInformation.of(CityPojo.class));FileSource<CityPojo>source=FileSource.forRecordStreamFormat(csvFormat,Path.fromLocalFile(...)).build();

对应 CSV:

Berlin|52.5167|13.3833|Germany|DE|Berlin|primary|3644826 San Francisco|37.7562|-122.443|United States|US|California||3592294 Beijing|39.905|116.3914|China|CN|Beijing|primary|19433000

更复杂的类型也能做(比如数组列),通过CsvSchema.ColumnType.ARRAY并指定数组元素分隔符:

CsvReaderFormat<ComplexPojo>csvFormat=CsvReaderFormat.forSchema(CsvSchema.builder().addColumn(newCsvSchema.Column(0,"id",CsvSchema.ColumnType.NUMBER)).addColumn(newCsvSchema.Column(4,"array",CsvSchema.ColumnType.ARRAY).withArrayElementSeparator("#")).build(),TypeInformation.of(ComplexPojo.class));

4、PyFlink:手动定义 CSV Schema(输出为 Row)

PyFlink 里通常自己建 schema,每一列映射为 Row 字段:

frompyflink.common.watermark_strategyimportWatermarkStrategyfrompyflink.tableimportDataTypesfrompyflink.datastreamimportStreamExecutionEnvironmentfrompyflink.datastream.connectors.file_systemimportFileSourcefrompyflink.formats.csvimportCsvReaderFormat,CsvSchema# 具体 import 以你环境包结构为准env=StreamExecutionEnvironment.get_execution_environment()schema=CsvSchema.builder()\.add_number_column('id',number_type=DataTypes.BIGINT())\.add_array_column('array',separator='#',element_type=DataTypes.INT())\.set_column_separator(',')\.build()source=FileSource.for_record_stream_format(CsvReaderFormat.for_schema(schema),CSV_FILE_PATH).build()ds=env.from_source(source,WatermarkStrategy.no_watermarks(),'csv-source')# ds 的 record 类型是 Row(具名字段 + 复合类型)# Types.ROW_NAMED(['id', 'array'], [Types.LONG(), Types.LIST(Types.INT())])

对应 CSV:

0,1#2#3 1, 2,1

补一个实战提醒:如果某列可能为空(比如上面的array),你后续算子处理时要把None/空数组的分支写好,否则很容易在 map/flat_map 里触发类型错误。

5、PyFlink:写 CSV(Bulk Format)

写 CSV 通常用CsvBulkWriters生成 BulkWriterFactory,再配合FileSink.for_bulk_format(...)

frompyflink.tableimportDataTypesfrompyflink.datastream.connectors.file_systemimportFileSinkfrompyflink.formats.csvimportCsvBulkWriters,CsvSchema# 具体 import 以你环境包结构为准schema=CsvSchema.builder()\.add_number_column('id',number_type=DataTypes.BIGINT())\.add_array_column('array',separator='#',element_type=DataTypes.INT())\.set_column_separator(',')\.build()sink=FileSink.for_bulk_format(OUTPUT_DIR,CsvBulkWriters.for_schema(schema)).build()ds.sink_to(sink)

6、运行模式:Batch / Streaming 都可用

CsvReaderFormat类似TextLineInputFormat,既可用于批也可用于流(持续监控目录等),具体取决于你用的 Source/RuntimeMode 以及文件系统是否支持持续发现。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/1 13:13:47

程序员必学:大模型构建四阶段详解,收藏这篇就够了

本文详细介绍了从零构建大语言模型的四个阶段&#xff1a;预训练&#xff08;通过海量语料教授基础知识&#xff09;、指令微调&#xff08;使模型遵循指令&#xff09;、偏好微调&#xff08;利用人类偏好数据通过RLHF对齐价值观&#xff09;和推理微调&#xff08;激发模型推…

作者头像 李华
网站建设 2026/7/24 5:59:39

HG_REPMGR autofailvoer自动故障转移

文章目录 文档用途详细信息 文档用途 HG_REPMGR自动故障转移配置参考 详细信息 配置集群自动故障转移&#xff08;failover&#xff09;&#xff0c;需要为集群中的每个节点开启 repmgrd 守护进程。当主节点出现故障后&#xff0c;会自动将合适的备节点提升为新主节点&#…

作者头像 李华
网站建设 2026/7/27 17:53:35

MySQL JOIN语法深度解析:从理论到实践的完整指南

目录 一、JOIN的本质与数学基础 二、内连接&#xff08;INNER JOIN&#xff09;的深层机制 三、外连接的完整语义解析 四、特殊连接类型的适用场景 五、JOIN性能优化的核心原则 六、JOIN与事务处理的交互影响 七、高级JOIN技术的实践应用 八、JOIN设计的最佳实践 结语 …

作者头像 李华
网站建设 2026/8/1 9:47:54

全网最全8个AI论文写作软件,助你轻松搞定本科毕业论文!

全网最全8个AI论文写作软件&#xff0c;助你轻松搞定本科毕业论文&#xff01; AI 工具&#xff0c;让论文写作不再焦虑 在当前的学术环境中&#xff0c;越来越多的本科生开始借助 AI 工具来辅助论文写作。这些工具不仅能够帮助学生快速生成内容&#xff0c;还能有效降低 AIGC&…

作者头像 李华
网站建设 2026/7/24 5:03:23

兽医AI推理TensorRT延迟砍半

&#x1f4dd; 博客主页&#xff1a;Jax的CSDN主页 兽医AI的“快”时代&#xff1a;TensorRT如何让动物诊断推理延迟砍半目录兽医AI的“快”时代&#xff1a;TensorRT如何让动物诊断推理延迟砍半 引言&#xff1a;兽医AI的延迟困境与破局点 一、兽医场景的特殊需求&#xff1a;…

作者头像 李华