ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache Beam 大数据实践:BigQueryIO.useBeamSchema 完全指南——用 Beam Schema 打通 BigQuery 读写

Apache Beam 大数据实践:BigQueryIO.useBeamSchema 完全指南——用 Beam Schema 打通 BigQuery 读写 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 的BigQueryIO是连接 Beam 流水线与 Google BigQuery 的核心 I/O 组件。useBeamSchema()方法决定了读写数据时采用哪套 schema 体系是 Beam 内部统一的 schema 表示还是 BigQuery 原生表 schema。本篇指南以 Beam 官方 Tour of Beam 教程beam-schema 单元为主线结合本仓库BigQueryIO源码实现讲解useBeamSchema的语义、底层调用链、完整可运行的 Java 示例以及使用限制帮助你写出类型安全、免手写转换函数的 BigQuery 读写流水线。一、useBeamSchema 是什么useBeamSchema是BigQueryIO提供的一个布尔开关用于指定在与 BigQuery 交互时使用哪种 schema 表示设为true调用useBeamSchema()Beam 使用其内部的 schema 表示即org.apache.beam.sdk.schemas.Schema。Beam 的 schema 体系支持更丰富的字段类型与更高级的 schema 操作嵌套结构、可空性、行转换函数等数据处理环节因此更灵活。设为false默认值Beam 使用 BigQuery 表的原生 schemaTableSchema。这适合希望写入格式与使用同一张表的其他工具保持兼容的场景。该开关默认关闭。从源码可见BigQueryIO.write()构建Write变换时显式执行.setUseBeamSchema(false)BigQueryIO.java只有显式调用useBeamSchema()才会切换为 Beam schema 模式。二、写入端useBeamSchema() 的完整用法2.1 官方示例代码Tour of Beam 教程给出的核心写法如下description.mdpipeline.apply(ReadFromBigQuery, BigQueryIO.write().to(mydataset.outputtable).useBeamSchema());其中BigQueryIO.write()创建一个Write变换把PCollection数据写入 BigQuery 表to(mydataset.outputtable)指定目标表名项目ID:数据集.表名或数据集.表名格式useBeamSchema()要求以PCollection元素的 schema 作为输出表的 schema不再需要手写TableSchema。2.2 方法签名与语义源码解读useBeamSchema()定义在 BigQueryIO.java/** * If true, then the BigQuery schema will be inferred from the input schema. If no * formatFunction is set, then BigQueryIO will automatically turn the input records into * TableRows that match the schema. */ public WriteT useBeamSchema() { return toBuilder().setUseBeamSchema(true).build(); }要点有二BigQuery 表 schema 从输入 PCollection 的 Beam schema 推断如果未设置withFormatFunctionBigQueryIO 会自动把输入记录转换为与 schema 匹配的TableRow省去手工new TableRow().set(...)的样板代码。2.3 底层实现开启后发生了什么在Write.expand的校验与装配阶段BigQueryIO.java开启useBeamSchema后按顺序执行if (getUseBeamSchema()) { checkArgument(input.hasSchema(), The input doesnt has a schema); // ① 输入必须带 schema optimizeWrites true; // ② 启用优化写入路径 checkArgument(avroRowWriterFactory null, avro avroFormatFunction is unsupported when using Beam schemas.); // ③ Avro 格式化互斥 if (formatFunction null) { // ④ 未提供 formatFunction 时自动生成 元素 → TableRow 转换 formatFunction TableRowFormatFunction.fromSerializableFunction( BigQueryUtils.toTableRow(input.getToRowFunction())); } // ⑤ 由 Beam schema 推断 BigQuery TableSchema TableSchema tableSchema BigQueryUtils.toTableSchema(input.getSchema()); dynamicDestinations new ConstantSchemaDestinations(...); }关键点① 输入必须有 schema若PCollection没有 schema流水线会在构建期直接抛出IllegalArgumentException(The input doesnt has a schema)。schema 来源可以是 POJO 上的DefaultSchema(JavaFieldSchema.class)注解、setRowSchema(...)或Create.of(...)携带的 schema。④ 自动转换转换函数由BigQueryUtils.toTableRow(input.getToRowFunction())生成见 BigQueryUtils.java内部把 BeamRow逐字段映射为TableRow。⑤ schema 推断BigQueryUtils.toTableSchema(Schema)BigQueryUtils.java把 Beam 字段类型一一转换为 BigQuery 字段类型最终通过ConstantSchemaDestinations绑定到目标表。2.4 可运行实战示例Tour of Beam 配套的完整示例位于 beam-schema/java-example/Task.java。下面是在其基础上整理出的完整可运行流水线含读取端到写入端// ① 定义 Beam Schema Schema inputSchema Schema.builder() .addField(id, Schema.FieldType.INT32) .addField(name, Schema.FieldType.STRING) .addField(age, Schema.FieldType.INT32) .build(); // ② 从 BigQuery 表读取默认以 TableRow 形式返回 PCollectionTableRow rows pipeline.apply( BigQueryIO.readTableRows().from(project-id.dataset.table)); // ③ 用 setRowSchema 为 PCollection 提供 Beam schema PCollectionObject typed rows .apply(MapElements.into(TypeDescriptor.of(Object.class)).via(it - it)) .setCoder(CustomCoder.of()) .setRowSchema(inputSchema); // ④ 写入另一张表useBeamSchema() 自动完成 schema 推断与类型转换 typed.apply(WriteToBigQuery, BigQueryIO.write() .to(mydataset.outputtable) .useBeamSchema()); pipeline.run();使用前需配置运行环境Task.javaSystem.setProperty(GOOGLE_APPLICATION_CREDENTIALS, /path/to/credential.json); PipelineOptions options PipelineOptionsFactory.fromArgs(args).create(); options.setTempLocation(gs://your-bucket); // 批量 load 路径需要临时 GCS 位置 options.as(BigQueryOptions.class).setProject(project-id); // 指定 GCP 项目2.5 与常用写入配置组合useBeamSchema()通常与Write的其他配置一起使用下表整理自源码枚举与默认值BigQueryIO.java、CreateDisposition、WriteDisposition配置方法可选值默认值说明useBeamSchema()true/falsefalse是否以 Beam schema 推断 BigQuery 表 schemawithCreateDisposition(...)CREATE_NEVER/CREATE_IF_NEEDEDCREATE_IF_NEEDED表不存在时是否自动创建CREATE_NEVER时表不存在则写入失败withWriteDisposition(...)WRITE_TRUNCATE/WRITE_APPEND/WRITE_EMPTYWRITE_EMPTYWRITE_TRUNCATE仅适用于FILE_LOADS方法对无界 PCollection 不支持withMethod(...)DEFAULT/FILE_LOADS/STREAMING_INSERTS/STORAGE_WRITE_API/STORAGE_API_AT_LEAST_ONCEDEFAULT有界输入默认走批量 load无界输入默认走流式插入withFormatFunction(...)自定义元素 → TableRow无与useBeamSchema()同时未设置时自动转换两者同时设置时以formatFunction为准withSchema(...)显式TableSchema无开启useBeamSchema后可省略典型组合示例批量覆盖写来自 BigQueryIO.javaPCollectionQuote quotes ...; // Quote 为带 DefaultSchema(JavaFieldSchema.class) 的 POJO quotes.apply(BigQueryIO.Quotewrite() .to(my-project:my_dataset.my_table) .useBeamSchema() .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_TRUNCATE));这里QuotePOJO 的 schema 由注解自动推断无需withSchema、withFormatFunction代码最简洁。三、读取端Beam Schema 如何接入useBeamSchema主要面向Write但读取端的 Beam schema 支持同样基于同一套 schema 转换机制。从源码看BigQueryIO.java当TypedRead同时设置了typeDescriptor、toBeamRowFn、fromBeamRowFn时Beam 会启用 schema 支持if (getTypeDescriptor() ! null getToBeamRowFn() ! null getFromBeamRowFn() ! null) { TableSchema tableSchema sourceDef.getTableSchema(bqOptions); // 获取目标表 schema ... beamSchema BigQueryUtils.fromTableSchema(tableSchema, builder.build()); // BigQuery → Beam }读取时通过withBeamRowConverters(...)BigQueryIO.java注册元素与Row的双向转换函数BigQuery 表 schema 会被转换为 BeamSchemaBigQueryUtils.fromTableSchema见 BigQueryUtils.java。也就是说写BeamSchema→TableSchema推断建表结构 自动Row → TableRow读BigQueryTableSchema→ BeamSchema配合withBeamRowConverters得到类型化PCollection。两方向互为镜像这正是 Beam schema 统一表示的体现。四、Beam Schema 写路径的两种落地方式开启useBeamSchema后根据所选写入方法数据会走两条不同的落地路径FILE_LOADS / STREAMING_INSERTS元素经自动生成的TableRowFormatFunction转成TableRowJSON 形式再交给批量 load 任务或流式插入 APIBigQueryIO.java。STORAGE_WRITE_API / STORAGE_API_AT_LEAST_ONCE使用StorageApiDynamicDestinationsBeamRow直接把 BeamRow翻译为 Storage Write API 的 proto 消息无需经过 JSON TableRow 往返转换性能开销更低BigQueryIO.javaif (getUseBeamSchema()) { checkArgument(!useSchemaUpdate, SchemaUpdateOptions are not supported when using Beam schemas); storageApiDynamicDestinations new StorageApiDynamicDestinationsBeamRow( dynamicDestinations, checkStateNotNull(elementSchema), checkStateNotNull(elementToRowFunction), getFormatRecordOnFailureFunction(), getRowMutationInformationFn() ! null); }同时开启useBeamSchema会自动设置optimizeWrites trueBigQueryIO.java与optimizedWrites()优化路径一致。五、使用限制与约束务必知晓源码中的checkArgument校验明确了以下硬性约束违反任一条件流水线都会在构建期失败输入 PCollection 必须携带 schema否则抛IllegalArgumentException(The input doesnt has a schema)BigQueryIO.java。不支持withAvroFormatFunction两者互斥报错信息为avroFormatFunction is unsupported when using Beam schemas。对应测试testWriteValidateFailsWithBeamSchemaAndAvroFormatFunctionBigQueryIOWriteTest.java。不支持自动 schema 更新Auto Schema UpdateAuto schema update not supported when using Beam schemas.BigQueryIO.java。STORAGE_WRITE_API 路径下不支持SchemaUpdateOptionsSchemaUpdateOptions are not supported when using Beam schemasBigQueryIO.java。类型映射并非一一对应源码注释明确说明BigQuery 的GEOGRAPHY等类型在 Beam schema 中没有精确对应物推断建表时需留意精度问题BigQueryIO.java。六、测试验证与底层佐证仓库中的单元测试直接印证了useBeamSchema的真实行为可作为实践参考批量写入 Beam schematestBatchSchemaWriteLoads用带SchemaCreate构造器的SchemaPojo元素写入配合withMethod(Method.FILE_LOADS)与useBeamSchema()最终断言TableRow内容与输入一致BigQueryIOWriteTest.java。流式插入 Beam schematestSchemaWriteStreams验证STREAMING_INSERTS路径下useBeamSchema()同样生效BigQueryIOWriteTest.java。动态目标表测试代码中useBeamSchema()可与按用户分表写入每用户一张表配合使用表明 Beam schema 模式同样支持动态目标表BigQueryIOWriteTest.java。此外在 Beam SQL / 声明式配置体系中该开关以use_beam_schema字段暴露BigQueryIOTranslation会把该字段写入 transform 配置并从配置中读回BigQueryIOTranslation.java说明useBeamSchema已纳入 Beam 的统一声明式 I/O 协议。七、总结BigQueryIO.useBeamSchema()让 Apache Beam 流水线以统一的 Beam Schema 体系与 BigQuery 交互写入时自动从输入PCollection推断表结构并生成转换函数读取时把 BigQuery 表 schema 反转为 BeamSchema。它大幅减少了手写TableRow转换与显式TableSchema的样板代码是构建类型安全、结构清晰的大数据管线的推荐做法。实践要点回顾写入前确保PCollection已带 schemaPOJO 注解或setRowSchema开启后不必再传withSchema与withFormatFunction两者之一仍显式提供时以显式为准注意与withAvroFormatFunction、自动 schema 更新、SchemaUpdateOptions的互斥限制配合withMethod、withCreateDisposition、withWriteDisposition可覆盖批量、流式、覆盖写等绝大多数场景。进一步学习可参考 Tour of Beam 的 big-query-io 单元目录内含 read-query、read-table、table-schema 等配套示例以及BigQueryIO的完整源码与测试BigQueryIO.java、BigQueryIOWriteTest.java。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam 大数据查询实战使用 BigQueryIO 连接器读写 Google BigQueryApache Beam 大数据查询实战使用 BigQueryIO 连接器读写 Google BigQuery Apache Beam 提供了开箱即用的 Big大数据批处理流处理数据工程Apache Beam 实战使用 BigQueryIO 向 Google BigQuery 写入数据的 Java 指南Apache Beam 实战使用 BigQueryIO 向 Google BigQuery 写入数据的 Java 指南 导读 本文以 Apache Beam大数据批处理流处理数据工程Apache Beam元数据管理Schema与类型安全实践Apache Beam元数据管理Schema与类型安全实践 在数据处理流程中元数据管理是确保数据一致性和处理效率的关键环节。Apache Beam作为统一批批处理流处理大数据上一篇Z3 TypeScript/WebAssembly 绑定z3-solver从 C API 自动生成到浏览器与 Node.js 的完整实践指南下一篇Relay 中 Disposable 类型详解dispose 语义、使用场景与源码实现创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表