如何用Java客户端库获取Google Dataflow项目内所有作业,求代码示例
列出Google Dataflow项目下所有作业的实现方案
使用Google Dataflow客户端库(Python示例)
首先安装依赖包:
pip install google-cloud-dataflow
代码示例:
from google.cloud import dataflow_v1beta3 def list_all_dataflow_jobs(project_id, region): # 初始化Dataflow客户端 client = dataflow_v1beta3.JobsV1Beta3Client() # 构造请求父路径 parent = f"projects/{project_id}/locations/{region}" # 分页获取所有作业 request = dataflow_v1beta3.ListJobsRequest(parent=parent) results = [] page_result = client.list_jobs(request=request) for page in page_result.pages: for job in page.jobs: results.append({ "job_id": job.id, "job_name": job.name, "state": job.state.name, "create_time": job.create_time }) return results # 调用示例 if __name__ == "__main__": PROJECT_ID = "your-project-id" REGION = "us-central1" # 替换为你的Dataflow作业所在区域 jobs = list_all_dataflow_jobs(PROJECT_ID, REGION) for job in jobs: print(f"作业ID: {job['job_id']}, 名称: {job['job_name']}, 状态: {job['state']}")
该示例会获取指定项目和区域下的全部Dataflow作业,返回包含作业核心信息的列表,同时处理了分页逻辑,确保不会遗漏作业。
使用Apache Beam Runner(Python示例)
如果使用Apache Beam的Dataflow Runner,可直接调用其内置的作业列表方法:
先安装Apache Beam的GCP依赖:
pip install apache-beam[gcp]
代码示例:
import apache_beam as beam from apache_beam.runners.dataflow.dataflow_runner import DataflowRunner def list_jobs_with_beam_runner(project_id, region): # 初始化Dataflow Runner实例 runner = DataflowRunner() # 列出指定项目和区域的所有作业 jobs = runner.list_jobs( project_id=project_id, region=region ) # 整理作业信息 job_list = [] for job in jobs: job_list.append({ "job_id": job.id, "job_name": job.name, "state": job.state, "created_at": job.created_at }) return job_list # 调用示例 if __name__ == "__main__": PROJECT_ID = "your-project-id" REGION = "us-central1" jobs = list_jobs_with_beam_runner(PROJECT_ID, REGION) for job in jobs: print(f"作业ID: {job['job_id']}, 名称: {job['job_name']}, 状态: {job['state']}")
此方法通过Beam Runner直接对接Dataflow API,返回的作业对象包含完整元数据,可根据需求提取更多字段。
内容的提问来源于stack exchange,提问作者VIKAS ROY
相关产品推荐
相关产品推荐

