You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.01 10:20:43