
大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 通过BigQueryIO这一连接器将 BigQuery 的 Storage Write API 能力无缝整合进统一批流编程模型让开发者可以用同一套代码完成「查询/表读取 → 转换 → 写入 BigQuery」的完整数据链路。本文以 Tour of Beam 学习模块 io/rest-api 中的 BigQuery API 实战示例为核心结合 Apache Beam 仓库内BigQueryIO的真实源码与示例工程系统讲解 Storage Write API 的核心概念、Java/Python 双语言读写姿势、动态目标表DynamicDestinations / table_side_inputs的两种实现方式以及写入 Disposition 等关键配置项的取舍。读完本文你将能够独立编写按业务字段动态分流写入多张 BigQuery 表的 Beam 管道并能读懂其背后的源码调用链。BigQuery Storage Write API统一流式与批量写入的单一入口description.md开篇即点明了本次技术的核心BigQuery Storage Write API 是一个面向 BigQuery 的统一数据摄取 API它把「流式摄取streaming ingestion」与「批量加载batch loading」合并为一个高性能 API。这意味着你可以用它在实时场景下把记录逐条流入 BigQuery在批处理场景下将任意规模的记录批量处理并在**单个原子操作single atomic operation**中一次性提交。从 Beam 的角度看BigQueryIO正是这一能力的封装入口。仓库源码 BigQueryIO.java 中public static Read read()是读取入口写入侧则由Write系列方法承担二者的组合构成了读 → 转换 → 写的完整管道。实战场景把天气观测数据按年份动态归档description.md中的 Java 示例给出了一个非常有代表性的业务场景通过 SQL 查询从apache-beam-testing.samples.weather_stations表读取 2007–2009 年的天气观测数据year、month、day、max_temperature使用自定义的WeatherDataPOJO 承载解析结果写入 BigQuery 时按照year字段动态分流为每一年生成一张独立的目标表如weather_2007、weather_2008、weather_2009。这个场景同时覆盖了 Storage Write API 的两大类用法从查询结果读取与按业务键动态写入多张表是学习BigQueryIO的最佳切入点。Java 读取BigQueryIO.read().fromQuery() 与类型化映射查询读取与记录解析Java 侧读取的核心是BigQueryIO.read()配合fromQuery(...)PCollectionWeatherData weatherData p.apply( BigQueryIO.read( (SchemaAndRecord elem) - { GenericRecord record elem.getRecord(); return new WeatherData( (Long) record.get(year), (Long) record.get(month), (Long) record.get(day), (Double) record.get(max_temperature)); }) .fromQuery( SELECT year, month, day, max_temperature FROM [apache-beam-testing.samples.weather_stations] WHERE year BETWEEN 2007 AND 2009) .withCoder(AvroCoder.of(WeatherData.class)));几个关键点值得展开BigQueryIO.read(SerializableFunctionSchemaAndRecord, T)这是类型化读取TypedRead的入口传入的函数负责把 BigQuery 返回的 AvroGenericRecord映射为业务对象。SchemaAndRecord同时携带 Avro 模式与记录本身函数内部通过elem.getRecord()拿到记录再按字段名取值。fromQuery(String)源码 BigQueryIO.java#L1147-L1153 中同时提供了fromQuery(String)与fromQuery(ValueProviderString)两个重载后者允许查询字符串在运行时由运行时参数如 Dataflow 模板参数动态提供适合模板化管道。from()与fromQuery()互斥源码 BigQueryIO.java#L1533 明确校验from() and fromQuery() are exclusive即直接读表与查询读表只能二选一。查询读取的临时数据集源码 BigQueryIO.java#L2353 说明使用fromQuery()时 Beam 会借助临时数据集保存查询结果若执行作业的服务账号缺少建数据集权限可以显式指定一个已有数据集对应withQueryLocation/ 临时数据集配置项。withCoder为下游PCollectionWeatherData指定 Coder。示例使用AvroCoder.of(WeatherData.class)保证分布式执行时对象序列化正确在 java-example/Task.java 中还能看到自定义CoderUser的完整写法encode/decode/verifyDeterministic供自定义类型参考。此外BigQueryIO.TypedRead源码 BigQueryIO.java#L2303-L2309还支持withCoder、withDeduplication等更多配置fromQuery()场景下可配合withQueryPriority()控制查询优先级from()与fromQuery()之外的方法互斥校验同样在源码中有所体现见 BigQueryIO.java#L1536。Java 写入DynamicDestinations 实现按年份分表读取到PCollectionWeatherData之后示例使用DynamicDestinations完成每一年一张表的动态写入weatherData.apply( BigQueryIO.WeatherDatawrite() .to( new DynamicDestinationsWeatherData, Long() { Override public Long getDestination(ValueInSingleWindowWeatherData elem) { return elem.getValue().year; } Override public TableDestination getTable(Long destination) { return new TableDestination( new TableReference() .setProjectId(writeProject) .setDatasetId(writeDataset) .setTableId(writeTable _ destination), Table for year destination); } Override public TableSchema getSchema(Long destination) { return new TableSchema() .setFields( ImmutableList.of( new TableFieldSchema() .setName(year) .setType(INTEGER) .setMode(REQUIRED), new TableFieldSchema() .setName(month) .setType(INTEGER) .setMode(REQUIRED), new TableFieldSchema() .setName(day) .setType(INTEGER) .setMode(REQUIRED), new TableFieldSchema() .setName(maxTemp) .setType(FLOAT) .setMode(NULLABLE))); } }) .withFormatFunction( (WeatherData elem) - new TableRow() .set(year, elem.year) .set(month, elem.month) .set(day, elem.day) .set(maxTemp, elem.maxTemp)) .withCreateDisposition(CreateDisposition.CREATE_IF_NEEDED) .withWriteDisposition(WriteDisposition.WRITE_TRUNCATE));DynamicDestinations 的三个核心回调DynamicDestinationsT, DestinationT是动态路由写入的抽象示例覆盖了它的全部三个方法方法职责本示例的实现getDestination(ValueInSingleWindowT)从每条元素中提取目标键决定数据去往哪张表返回elem.getValue().yeargetTable(DestinationT)根据目标键生成目标表的引用与描述拼接writeTable _ destination并附带描述文本getSchema(DestinationT)根据目标键返回目标表的 Schema声明 year/month/day/maxTemp 四个字段的类型与模式值得注意的细节Schema 与类型year、month、day为INTEGER且REQUIRED必填maxTemp为FLOAT且NULLABLE可空。BigQuery 类型名INTEGER/FLOAT与 Java 侧字段类型Long/Double由withFormatFunction负责桥接。getTable的投影字段示例中每个目标表都只保存了当年数据通过表名后缀区分getDestination的返回值即表名后缀体现按值分表的经典做法。withFormatFunction 与写入 DispositionwithFormatFunction把业务对象WeatherData转换为 BigQuery 的TableRow。这是与读取侧BigQueryIO.read(解析函数)对称的序列化桥。withCreateDisposition(CREATE_IF_NEEDED)目标表不存在时自动创建与getSchema配合使用这是表可以不存在的动态写入前提。withWriteDisposition(WRITE_TRUNCATE)每次运行时清空并重写目标表适合周期性全量归档如按年份重建天气归档表的场景。若需追加数据应改用WRITE_APPEND若希望表存在即报错则用WRITE_EMPTY。仓库中的同构示例学习目录 java-example/Task.java 提供了一份结构完全相同的可运行版本它从projectId.dataset.table读取用户数据id/name/age以id作为DynamicDestinations的目标键为每个用户生成table_id目标表示例还展示了PipelineOptions、setTempLocationgs://btestq、BigQueryOptions.setProject等运行期配置以及LogOutputT日志 DoFn 的写法可直接作为本地/Dataflow 运行的起点。同一单元 unit-info.yaml 标明该课程同时覆盖 Java 与 Python复杂度为 ADVANCED。Python 写入WriteToBigQuery 与 table_side_inputs 动态分流Python 侧description.md的{{if (eq .Sdk python)}}分支展示了用side input旁路输入实现动态目标表的另一条路线fictional_characters_view beam.pvalue.AsDict( pipeline | CreateCharacters beam.Create([(Yoda, True),(Obi Wan Kenobi, True)])) def table_fn(element, fictional_characters): if element in fictional_characters: return my_dataset.fictional_quotes else: return my_dataset.real_quotes quotes | WriteWithDynamicDestination beam.io.WriteToBigQuery( table_fn, schematable_schema, table_side_inputs(fictional_characters_view, ), write_dispositionbeam.io.BigQueryDisposition.WRITE_TRUNCATE, create_dispositionbeam.io.BigQueryDisposition.CREATE_IF_NEEDED)关键设计拆解table_fn作为 table 参数WriteToBigQuery的table参数既可以传固定字符串dataset.table也可以传一个函数。传入函数时Beam 会逐元素调用该函数用返回值决定每条记录写入哪张表——这就是 Python 版的动态目标表。table_side_inputs当table是函数且函数签名中带有额外的旁路参数时用table_side_inputs把这些PCollection以AsDict/AsList/AsSingleton等形式注入。示例把fictional_characters_view一个AsDict视图注入table_fn据此判断这条 quote 出自虚构角色还是真实角色从而分流到fictional_quotes或real_quotes两张表。Python SDK 源码 bigquery.py#L2177 对table_side_inputs的注释明确指出它是一个由AsSideInputPCollection 组成的 tuple会在调用 table 函数时被一并展开*self.table_side_inputs见 bigquery.py#L1974。schema参数目标表的 Schema 定义可以在 python-example/task.py 中看到用bigquery.TableSchema()/bigquery.TableFieldSchema()编程式构造 Schema 的完整写法其中字段模式nullable/required与 Java 侧的Mode一一对应。写入 Dispositionwrite_dispositionWRITE_TRUNCATE、create_dispositionCREATE_IF_NEEDED与 Java 侧语义一致Python SDK 在 bigquery.py#L584 定义了BigQueryDisposition枚举WRITE_EMPTY/WRITE_TRUNCATE/WRITE_APPEND/CREATE_IF_NEEDED/CREATE_NEVER两两组合可覆盖绝大多数写入语义。Java 与 Python 动态写入的对比与选型维度JavaDynamicDestinationsPythontable_fn table_side_inputs目标键来源元素本身getDestination(elem)元素 可选的旁路输入表名/描述getTable(destination)返回TableDestination含描述函数返回表名字符串SchemagetSchema(destination)按目标键返回统一schema参数或表级 schema外部参考数据需自行用 side input/state 结合原生table_side_inputs注入典型场景按元素字段年份、用户 id分表按是否命中字典等外部知识分表选型建议当分流逻辑仅依赖元素自身的字段如按年份、按用户 id时Java 的DynamicDestinations与 Python 的table_fn均合适当分流依赖额外的参考数据集如白名单、维度表时Python 的table_side_inputs是更直接的表达方式Java 侧则需要借助 side input 在管道中先行View化再传入DynamicDestinations。从学习模块到生产运行前提与限制说明本单元属于 io 模块IO Connectors复杂度 ADVANCED下的rest-api子单元与text-io、big-query-io、kafka-io并列。将上述代码投入生产前需要明确以下前提环境与凭证必须配置有效的 GCP 项目与凭证。Java 侧可通过BigQueryOptions指定项目并在运行环境中设置GOOGLE_APPLICATION_CREDENTIALS见 Task.java#L68-L71fromQuery()读取还需要作业账号具备创建临时数据集的权限否则需显式指定临时数据集。临时位置setTempLocation(gs://btestq)之类的 GCS 临时目录是 BigQueryIO 落中间结果查询结果暂存、写入暂存文件的必要配置示例中的 bucket 需替换为你自己的可写位置。原子性语义Storage Write API 的批量加载 单次原子提交特性由 BigQuery 侧保证Beam 管道侧的WRITE_TRUNCATE等 Disposition 决定的是重跑管道时旧数据如何处理二者需结合业务对幂等性的要求一起设计。版本约束示例依赖 Beam Java / Python SDK 中apache_beam.io.gcp.bigquery与org.apache.beam.sdk.io.gcp.bigquery包Python 示例位于 python-example/task.py具体行为以你使用的 Beam 版本为准本文所述源码位置均指向当前仓库。小结围绕 io/rest-api 描述文档 的 BigQuery API 主题本文完整覆盖了Storage Write API流批统一、原子提交的核心定位Java 侧BigQueryIO.read().fromQuery()的查询读取与DynamicDestinations按年分表写入Python 侧WriteToBigQuery结合table_side_inputs的字典驱动分流以及CREATE_IF_NEEDED/WRITE_TRUNCATE等写入 Disposition 的语义与选型。从 BigQueryIO.java 的from()/fromQuery()互斥校验到 bigquery.py 的table_side_inputs展开逻辑仓库源码印证了文档示例背后的真实调用链。以此为模板你可以轻松把按任意业务键动态分表的模式推广到自己的批流管道中。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam 实践BigQuery 表模式Table Schema与 DynamicDestinations 动态目标写入Apache Beam 实践BigQuery 表模式Table Schema与 DynamicDestinations 动态目标写入 导读 本指南以 Ap大数据批处理流处理数据工程Apache Beam 实战使用 BigQueryIO 向 Google BigQuery 写入数据的 Java 指南Apache Beam 实战使用 BigQueryIO 向 Google BigQuery 写入数据的 Java 指南 导读 本文以 Apache Beam大数据批处理流处理数据工程Apache Beam BigQueryIO 读取 BigQuery 表readTableRows实战指南Apache Beam BigQueryIO 读取 BigQuery 表readTableRows实战指南 Apache Beam 的 BigQueryIO大数据批处理流处理数据工程上一篇炉石传说HsMod插件终极游戏体验增强完整指南下一篇深度解析HsMod炉石传说插件5大核心技术模块与安全部署实战创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考