如何在Kubeflow流水线中获取ModelBatchPredictOp输出的big_query_table的tableId
Kubeflow流水线获取ModelBatchPredictOp输出的BigQuery表属性
问题背景
我在Kubeflow流水线中使用ModelBatchPredictOp组件,它可生成batchpredictionjob、big_query_table和gcs_output_directory三个工件,流水线运行正常。需要获取big_query_table工件的tableId属性用于后续BigqueryQueryJobOp组件,同时从其URI中提取完整表路径。
以下是相关组件代码:
# 批量预测组件 batch_predict_op = ModelBatchPredictOp( project=project_id, location=DEFAULT_VERTEX_REGION, instances_format = 'bigquery', predictions_format = 'bigquery', model=importer_spec.outputs['artifact'], job_display_name='teste_batch_predict', bigquery_source_input_uri=f'bq://{input_data_table_ref}', bigquery_destination_output_uri= f'bq://{output_bq}', ).after(input_data_table_op) top_predictions_table_ref = f'{project_id}.{bigquery_dataset}.test' # 基于前序组件创建表的组件 top_predictions_op = bq.BigqueryQueryJobOp( project_id, location = bigquery_job_location, query = predict_dataset.get_query( output_table = top_predictions_table_ref, source_table = batch_predict_op.outputs['bigquery_output_table'], query_name = 'query_top_100.sql', DEBUG = DEBUG) ).after(batch_predict_op)
解决方案
ModelBatchPredictOp输出的big_query_table是标准Artifact对象,可通过其metadata和uri字段直接获取所需信息:
1. 获取tableId属性
直接从Artifact的metadata字典中提取tableId字段:
prediction_table_id = batch_predict_op.outputs['big_query_table'].metadata['tableId']
2. 获取完整表路径
Artifact的uri字段格式为bq://project.dataset.table,移除前缀即可得到可用的完整表路径:
prediction_full_table_path = batch_predict_op.outputs['big_query_table'].uri.replace('bq://', '')
修改后的完整代码
# 批量预测组件 batch_predict_op = ModelBatchPredictOp( project=project_id, location=DEFAULT_VERTEX_REGION, instances_format='bigquery', predictions_format='bigquery', model=importer_spec.outputs['artifact'], job_display_name='teste_batch_predict', bigquery_source_input_uri=f'bq://{input_data_table_ref}', bigquery_destination_output_uri=f'bq://{output_bq}', ).after(input_data_table_op) top_predictions_table_ref = f'{project_id}.{bigquery_dataset}.test' # 提取BigQuery表的tableId和完整路径 prediction_table_id = batch_predict_op.outputs['big_query_table'].metadata['tableId'] prediction_full_table_path = batch_predict_op.outputs['big_query_table'].uri.replace('bq://', '') # 后续消费组件(根据SQL需求选择使用tableId或完整路径) top_predictions_op = bq.BigqueryQueryJobOp( project_id, location=bigquery_job_location, query=predict_dataset.get_query( output_table=top_predictions_table_ref, source_table=prediction_full_table_path, # 若SQL仅需表名则替换为prediction_table_id query_name='query_top_100.sql', DEBUG=DEBUG ) ).after(batch_predict_op)
注意事项
- 确认组件输出字段名:部分Kubeflow版本中,
big_query_table可能被命名为bigquery_output_table,需根据组件定义调整输出引用。 - URI格式兼容性:
bq://前缀是Vertex AI Batch Prediction输出BigQuery表URI的标准格式,若后续格式变动,需调整前缀移除逻辑。
内容的提问来源于stack exchange,提问作者Openworld
相关产品推荐
相关产品推荐

