
批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载Apache Beam 的CsvIO连接器位于org.apache.beam.sdk.io.csv包源码见 CsvIO.java为 Java SDK 提供了基于 Schema 的 CSV 文件读写能力。本文以代码生成实战场景为主线完整讲解如何用CsvIO.Write将一个PCollectionPOJO 或Row写出为带表头、支持自定义注释与分片的 CSV 文件并结合仓库源码剖析其底层实现机制、字段类型约束与高级配置项帮助你写出可直接运行、可深入调优的生产级代码。阅读本文后你将掌握基于DefaultSchema的 POJO 建模、PipelineOptions命令行参数解析、CSVFormat的定制表头、注释标记、字段子集、withNumShards/withSuffix/withCompression等写出控制以及 CsvIO 内部Schema → Row → CSV 字符串 → TextIO 落盘的完整调用链。一、CsvIO 是什么位于 Beam 官方 SDK 的 CSV 连接器CsvIO是 Apache Beam Java SDK 中专门负责 CSV 格式读写的连接器其核心实现位于仓库的sdks/java/io/csv模块。它依赖 Apache Commons CSV 库org.apache.commons.csv.CSVFormat来完成 CSV 文本层的格式化与解析而文件读写层则复用 Beam 自带的TextIO/FileBasedSink能力因此天然继承了 Beam 文件系统抽象可写到本地、GCS、S3 等任意 Beam 支持的文件系统。从类结构看见 CsvIO.javaCsvIO对外暴露的核心入口包括静态方法用途CsvIO.write(to, csvFormat)写出自定义 Java 类型POJO/AutoValue的PCollectionSchema 由 Beam 自动推断CsvIO.writeRows(to, csvFormat)写出PCollectionRowCsvIO.parse(klass, csvFormat)将 CSV 字符串记录解析为自定义类型读取方向CsvIO.parseRows(schema, csvFormat)将 CSV 字符串记录解析为Row读取方向本文聚焦写出方向CsvIO.Write读取方向的parse/parseRows在文中作为延伸补充。二、核心示例用 CsvIO 将内存数据写出为 CSV 文件下面的完整示例来自仓库的代码生成文档 08_io_csv.md它演示了一个最小的端到端流程定义带 Schema 的 POJO → 用Create构造PCollection→ 用CsvIO.write()写出到指定路径并限制只生成 1 个分片文件。import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.io.csv.CsvIO; import org.apache.beam.sdk.options.Description; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.schemas.JavaFieldSchema; import org.apache.beam.sdk.schemas.annotations.DefaultSchema; import org.apache.beam.sdk.transforms.Create; import org.apache.beam.sdk.values.PCollection; import org.apache.commons.csv.CSVFormat; import java.io.Serializable; import java.util.Arrays; import java.util.List; public class WriteCsvFile { // ExampleRecord is a POJO that represents the data to be written to the CSV file DefaultSchema(JavaFieldSchema.class) public static class ExampleRecord implements Serializable { public int id; public String month; public String amount; public ExampleRecord() { } public ExampleRecord(int id, String month, String amount) { this.id id; this.month month; this.amount amount; } } public interface WriteCsvFileOptions extends PipelineOptions { Description(A file path to write CSV files to) Validation.Required String getFilePath(); void setFilePath(String filePath); } public static void main(String[] args) { WriteCsvFileOptions options PipelineOptionsFactory.fromArgs(args) .withValidation().as(WriteCsvFileOptions.class); Pipeline p Pipeline.create(options); ListExampleRecord rows Arrays.asList( new ExampleRecord(1, January, $1000), new ExampleRecord(2, February, $2000), new ExampleRecord(3, March, $3000)); CSVFormat csvFormat CSVFormat.DEFAULT.withHeaderComments(CSV file created by Apache Beam) .withCommentMarker(#); p.apply(Create collection, Create.of(rows)) .apply( Write to CSV file, CsvIO.ExampleRecordwrite() .to(options.getFilePath()) .withNumShards(1)); p.run(); } }运行方式文件路径通过命令行参数注入# 本地 DirectRunner 运行 ./gradlew -p sdks/java/io/csv run \ --args--filePath/tmp/beam-csv-output # 或直接指定任意 Beam Runner 支持的文件系统路径例如 GCS # --filePathgs://your-bucket/path/to/prefix由于示例中CSVFormat.DEFAULT未显式指定表头且CsvIO.Write会在表头缺失时依据 Schema 自动生成最终输出的000000000001-of-000000000001.csv前缀为--filePath指定值内容大致为# CSV file created by Apache Beam amount,id,month $1000,1,January $2000,2,February $3000,3,March注意两点字段顺序按字段名排序amount,id,month且第一行是使用#注释标记书写的表头注释——这正是withCommentMarker(#)的作用。三、示例逐步拆解3.1 用 DefaultSchema 为 POJO 声明 SchemaDefaultSchema(JavaFieldSchema.class) public static class ExampleRecord implements Serializable { public int id; public String month; public String amount; }CsvIO.Write要求输入PCollection必须携带 Schema详见下文源码剖析。对于普通 POJO最简单的方式就是加上DefaultSchema(JavaFieldSchema.class)注解让 Beam 在运行时依据公有字段自动推断Schema。JavaFieldSchema属于org.apache.beam.sdk.schemas包除它之外Beam 还支持JavaBeanSchema基于 getter/setter 的 JavaBean 规范与AutoValueSchema配合 Google AutoValue详见源码 CsvIO.java 中Transaction类的示例。3.2 用 PipelineOptions 解析命令行参数public interface WriteCsvFileOptions extends PipelineOptions { Description(A file path to write CSV files to) Validation.Required String getFilePath(); void setFilePath(String filePath); }Beam 的PipelineOptions是一种通过接口方法自动生成实现的配置模式接口中的getXxx()/setXxx()方法对对应一个同名命令行参数本例为--filePath...。Description用于生成帮助信息Validation.Required来自org.apache.beam.sdk.options.Validation声明该参数为必填未提供时PipelineOptionsFactory.fromArgs(args).withValidation()会在启动阶段直接报错避免管道跑到一半才发现路径缺失。3.3 定制 CSVFormat表头注释与注释标记CSVFormat csvFormat CSVFormat.DEFAULT.withHeaderComments(CSV file created by Apache Beam) .withCommentMarker(#);这里演示了两个 Apache Commons CSV 的配置项withHeaderComments(...)在文件首部输出若干行注释常用于记录报表标题、生成日期或操作者信息withCommentMarker(#)指定注释行的前缀字符。二者必须成对使用CsvIO 在构建Write时会对配置做校验见 CsvIO.java如果设置了withHeaderComments而未设置withCommentMarker会抛出IllegalArgumentException提示CSVFormat withCommentMarker required when withHeaderComments。3.4 组装管道Create CsvIO.writep.apply(Create collection, Create.of(rows)) .apply( Write to CSV file, CsvIO.ExampleRecordwrite() .to(options.getFilePath()) .withNumShards(1)); p.run();Create.of(rows)将内存中的ListExampleRecord变成PCollectionExampleRecord并从泛型推断其 SchemaCsvIO.ExampleRecordwrite()返回CsvIO.WriteExampleRecord.to(path)指定输出文件前缀目录 文件名前缀实际文件会追加分片序号与.csv后缀.withNumShards(1)固定只生成 1 个分片文件适合小数据集或期望单文件的场景。四、深入源码CsvIO.Write 的底层工作原理CsvIO.Write是一个PTransformPCollectionT, WriteFilesResultString其核心逻辑在expand方法中见 CsvIO.java整个写出过程可以概括为四个阶段校验输入必须携带 Schemaexpand首先检查input.hasSchema()不满足则抛出IllegalArgumentException错误信息明确指出CsvIO requires an input Schema. Note that only Row or user classes are supported. Consider using TextIO or FileIO directly when writing primitive types。也就是说CsvIO.Write只接受Row或带 Schema 注解的用户类不支持PCollectionString、PCollectionInteger这类原始类型——这种情况请改用TextIO或FileIO。转成 Row通过MapElements和输入自带的行转换函数getToRowFunction()将T统一映射为PCollectionRow并绑定RowCoder。Row → CSV 字符串用CsvRowConversions.RowToCsv见 CsvRowConversions.java把每个Row按表头顺序取值后交给CSVFormat.format(values)序列化为一行 CSV 文本。TextIO 落盘PCollectionString最后交给底层TextIO.Write带withOutputFilenames()写出为文件。4.1 表头自动生成默认按字段名排序如果调用时没有通过CSVFormat.withHeader(...)显式指定表头CsvIO 会调用buildHeaderFromSchemaIfNeededCsvIO.java从 Schema 自动生成csvFormat csvFormat.withHeader(schema.sorted().getFieldNames().toArray(new String[0]));schema.sorted()意味着默认输出顺序是字段名的字典序而不是 POJO 中字段的声明顺序。这正是前文输出中顺序为amount,id,month的原因。若需要控制列的顺序与子集必须显式使用withHeader。4.2 表头与注释的写入细节writeWithCSVFormatHeaderAndCommentsCsvIO.java做了三件事对每个headerComments元素拼成commentMarker comment一行例如# CSV file created by Apache Beam用withSkipHeaderRecord()防止 Commons CSV 在格式化时输出两份表头将注释 表头行整体作为TextIO.Write的withHeader(...)传入因此每个分片文件都会重复写入同样的表头与注释。4.3 底层文件写出能力从何而来createDefaultTextIOWriteCsvIO.java显示CsvIO 默认构造的是TextIO.write().to(to).withSuffix(DEFAULT_FILENAME_SUFFIX); // DEFAULT_FILENAME_SUFFIX .csv因此CsvIO.Write的很多方法与TextIO.Write一一对应、语义一致如withNumShards、withSuffix、withoutSharding等默认文件后缀为.csv。这也是Write的返回类型是WriteFilesResultString包含getPerDestinationOutputFilenames()等信息的原因。五、支持的字段类型与 Schema 约束CSV 是扁平化文本格式无法表达嵌套结构。因此CsvIO.Write只支持不含嵌套类型的 Schema 字段合法集合由VALID_FIELD_TYPE_SET定义见 CsvIO.java合法 FieldType对应 Java 类型示例BYTEbyte/ByteBOOLEANboolean/BooleanDATETIMEInstant等时间类型DECIMALBigDecimalDOUBLEdouble/DoubleINT16short/ShortINT32int/IntegerINT64long/LongFLOATfloat/FloatSTRINGString不支持的嵌套类型包括ROW嵌套行与ARRAY重复字段。设计 POJO 时请保持字段为上述标量类型嵌套结构应先通过FlatMap等变换展平再交给 CsvIO 写出。六、高级配置完整控制输出文件CsvIO.Write在 CsvIO.java 中提供了与TextIO.Write对齐的完整配置链按需组合使用方法作用典型场景withNumShards(Integer)固定每个窗口的分片数量小数据量、希望得到确定数量文件withoutSharding()强制输出单个文件空分片模板导出单文件供人工查看withShardTemplate(String)自定义分片文件名模板ShardNameTemplate控制分片命名格式withSuffix(String)覆盖默认.csv后缀输出.tsv等变体配合自定义 delimiterwithCompression(Compression)指定输出压缩方式Compression.GZIP等压缩存储withWindowedWrites()按输入元素的窗口分别写文件流式管道按窗口输出withTempDirectory(ResourceId)指定临时文件目录跨文件系统写临时文件withNoSpilling()禁止数据溢出到磁盘WriteFiles.withNoSpilling内存充足、追求速度withWritableByteChannelFactory(...)自定义写通道工厂接入自定义编码/加密通道例如期望输出 GZip 压缩的单文件p.apply(Create.of(rows)) .apply(CsvIO.ExampleRecordwrite() .to(options.getFilePath()) .withoutSharding() .withCompression(Compression.GZIP));七、控制列顺序与子集withHeader 的正确用法要打破默认的字典序、只写部分字段使用CSVFormat.withHeader(...)显式声明例如输出transactionId,purchaseAmount两列p.apply(transactions) .apply(CsvIO.Transactionwrite( path/to/folder/prefix, CSVFormat.DEFAULT.withHeader(transactionId, purchaseAmount)));使用withHeader时需注意 CsvIO.java 列出的三条约束每个表头列名必须与 Schema 字段名匹配且大小写敏感Matching is case sensitive匹配上的字段必须是VALID_FIELD_TYPE_SET中的合法类型只有在withAllowDuplicateHeaderNames(true)时才允许表头列名重复。八、写出 PCollection writeRows 与 writeRowsTo如果你的管道数据本身就是 Schema 化的Row例如从BigQueryIO、AvroIO读取后未转换则不必定义 POJO直接用writeRows/writeRowsTo。仓库源码 CsvIO.java 给出了一个Transaction示例先用DefaultSchemaProvider从类型推导出Schema再用Row.withSchema(schema).withFieldValue(...)构造数据行DefaultSchemaProvider defaultSchemaProvider new DefaultSchemaProvider(); Schema schema defaultSchemaProvider.schemaFor(TypeDescriptor.of(Transaction.class)); PCollectionRow transactions pipeline.apply(Create.of( Row.withSchema(schema).withFieldValue(bank, A) .withFieldValue(purchaseAmount, 10.23) .withFieldValue(transactionId, 12345).build(), Row.withSchema(schema).withFieldValue(bank, B) .withFieldValue(purchaseAmount, 54.65) .withFieldValue(transactionId, 54321).build())); transactions.apply( CsvIO.writeRowsTo(gs://bucket/path/to/folder/prefix, CSVFormat.DEFAULT));writeRowsTo等价于writeRows(...).to(...)的便捷调用产出同样带表头的 CSV 文件。九、不支持的 CSVFormat 属性为保证写出行为可控CsvIO.Write明确不支持以下CSVFormat属性一旦启用会抛出IllegalArgumentException见 CsvIO.javawithAllowMissingColumnNameswithAutoFlushwithIgnoreHeaderCasewithIgnoreSurroundingSpaces其中withIgnoreHeaderCase与表头匹配的大小写敏感约束相呼应其余属性与单文件写语义或解析行为冲突。设计管道时请绕开这些配置。十、延伸CsvIO 的读取能力与测试佐证虽然本文主题是写出但 CsvIO 也提供对称的解析能力可作为完整 CSV 管道的拼图CsvIO.parse(SomeDataClass.class, csvFormat)把TextIO读出的 CSV 字符串记录解析为自定义类型结果CsvIOParseResultT同时携带getOutput()成功记录与getErrors()解析失败记录可接死信队列CsvIO.parseRows(schema, csvFormat)解析为Row。仓库测试 CsvIOTest.java 验证了关键行为解析带注释、带引号内换行foo\nbar、带内嵌逗号foo$,bar的记录以及非法 CSVFormat 抛IllegalArgumentExceptionSchema 与 CSVFormat 不匹配抛异常无 Schema 注解的类抛IllegalStateException等边界情况可作为你编写相似管道的回归参考。十一、参考资源本文关联的代码生成文档08_io_csv.mdCsvIO 核心实现CsvIO.javaRow 与 CSV 互转实现CsvRowConversions.java单元测试CsvIOTest.java同类 I/O 代码生成文档Kafka 等learning/prompts/code-generation/java总结CsvIO.Write让把 Beam 管道结果落成 CSV 文件这件事变成三行配置——定义带 Schema 的 POJO或Row、按需定制CSVFormat、调用CsvIO.write().to(...)。理解其背后Schema 校验 → Row 转换 → 表头生成 → TextIO 分片写出的实现链路后你就能准确预判字段排序、表头注释、分片与压缩等每一个细节写出既符合 Beam 惯例又可稳定运行于批/流场景的 CSV 输出管道。赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐Apache Beam Java CsvIO 实战使用 CsvIO 将 Schema 化数据写入 CSV 文件与源码级解析Apache Beam Java CsvIO 实战使用 CsvIO 将 Schema 化数据写入 CSV 文件与源码级解析 Apache Beam 的 Csv大数据批处理流处理数据工程Apache Beam 使用 CsvIO 写入 CSV 文件完整实战指南与源码解析Apache Beam 使用 CsvIO 写入 CSV 文件完整实战指南与源码解析 本篇技术指南以 Apache Beam Java SDK 的 CsvIO大数据批处理流处理数据工程Apache Beam Java 使用 CsvIO 写入 CSV 文件从 Schema 推断到分片输出的完整实战解析Apache Beam Java 使用 CsvIO 写入 CSV 文件从 Schema 推断到分片输出的完整实战解析 Apache Beam 的 Java S批处理流处理大数据上一篇Apache Pulsar 2.2.0 命令行工具完全指南pulsar、pulsar-client、pulsar-perf 与配套工具详解下一篇解读 docker-selenium 4.29.0 Firefox 129 镜像发布记录从打标签脚本看 node-firefox 镜像的版本命名体系创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考