PHP BigQuery Write Client批量写入:如何添加实际表数据?
使用BigQuery Write API Pending流批量写入数据(PHP)
核心修改说明
你需要在构建AppendRowsRequest阶段添加序列化后的表数据,同时补充流提交的收尾步骤,以下是完整的可运行代码示例:
use Google\Cloud\BigQuery\BigQueryClient; use Google\Cloud\BigQuery\Storage\V1\Client\BigQueryWriteClient; use Google\Cloud\BigQuery\Storage\V1\CreateWriteStreamRequest; use Google\Cloud\BigQuery\Storage\V1\WriteStream; use Google\Cloud\BigQuery\Storage\V1\AppendRowsRequest; use Google\Cloud\BigQuery\Storage\V1\ProtoData; use Google\Cloud\BigQuery\Storage\V1\Row; use Google\Cloud\BigQuery\Storage\V1\SubmitWriteStreamRequest; // 初始化客户端 $client = new BigQueryClient(); $writeClient = new BigQueryWriteClient(); // 目标表配置 $datasetId = '你的数据集ID'; $tableId = '你的表ID'; $table = $client->dataset($datasetId)->table($tableId); $tableIdentity = $table->identity(); // 构建表资源路径 $parent = $writeClient::tableName( project: $tableIdentity['projectId'], dataset: $tableIdentity['datasetId'], table: $tableIdentity['tableId'] ); // 创建Pending类型写入流 $writeStreamOptions = ['type' => WriteStream\Type::PENDING]; $streamCreationRequest = CreateWriteStreamRequest::build( $parent, new WriteStream($writeStreamOptions) ); $writeStream = $writeClient->createWriteStream($streamCreationRequest); $streamName = $writeStream->getName(); // -------------------------- // 1. 准备批量写入数据 // -------------------------- // 数据结构需与目标表Schema完全匹配 $rowsData = [ ['id' => 1, 'username' => 'Alice', 'age' => 30], ['id' => 2, 'username' => 'Bob', 'age' => 25], // 可添加更多数据行 ]; // -------------------------- // 2. 将数据序列化为BigQuery兼容格式 // -------------------------- $tableSchema = $table->schema(); $protoRows = []; foreach ($rowsData as $rowData) { $row = new Row(); foreach ($tableSchema->fields() as $field) { $fieldName = $field->name(); $value = $rowData[$fieldName] ?? null; // 根据字段类型设置对应值(需适配所有用到的类型) switch ($field->type()) { case 'INTEGER': $row->setInt64Value($value); break; case 'STRING': $row->setStringValue($value); break; case 'FLOAT': $row->setFloat64Value($value); break; // 其他类型如BOOLEAN、DATE等需对应处理 } $row->addValues($row->getInt64Value() ?? $row->getStringValue() ?? $row->getFloat64Value() ?? null); } $protoRows[] = $row; } // 封装序列化后的数据 $protoData = new ProtoData(); $protoData->setRows(new \Google\Cloud\BigQuery\Storage\V1\RowBatch(['rows' => $protoRows])); // -------------------------- // 3. 构建带数据的AppendRowsRequest // -------------------------- $appendRowRequest = AppendRowsRequest::build( $streamName, ['protoRows' => $protoData] ); // -------------------------- // 4. 写入数据到双向流 // -------------------------- $bidiStream = $writeClient->appendRows(); $bidiStream->write($appendRowRequest); // 验证写入结果 $response = $bidiStream->read(); if ($response->getAppendResult()->getError()) { throw new \Exception('写入失败: ' . $response->getAppendResult()->getError()->getMessage()); } // -------------------------- // 5. 提交Pending流(必须执行,否则数据不会生效) // -------------------------- $submitRequest = SubmitWriteStreamRequest::build($streamName); $writeClient->submitWriteStream($submitRequest); echo "批量数据写入完成!";
关键注意事项
- Schema一致性:写入数据的字段名、类型、顺序必须与目标表Schema完全匹配,否则会写入失败。
- 类型适配:不同字段类型需调用对应的
setXXXValue方法赋值,避免类型不兼容问题。 - 流提交:Pending类型流必须调用
submitWriteStream,才能将临时写入的数据最终合并到目标表中。 - 分批处理:若数据量过大,建议拆分多个批次写入,避免单次请求超时。
内容的提问来源于stack exchange,提问作者Russell
相关产品推荐
相关产品推荐

