【问题标题】:How to specify which GCP project to use when triggering a pipeline through Data Fusion operator on Cloud Composer如何在 Cloud Composer 上通过 Data Fusion 算子触发管道时指定要使用的 GCP 项目
【发布时间】:2021-12-09 07:38:35
【问题描述】:

我需要通过 DAG 内的数据融合运算符 (CloudDataFusionStartPipelineOperator) 触发位于名为 myDataFusionProject 的 GCP 项目上的数据融合管道,该 DAG 的 Cloud Composer 实例位于另一个名为 myCloudComposerProject 的项目上。

我使用official documentationsource code 编写了大致类似于以下sn-p 的代码:

LOCATION = "someLocation"
PIPELINE_NAME = "myDataFusionPipeline"
INSTANCE_NAME = "myDataFusionInstance"
RUNTIME_ARGS = {"output.instance":"someOutputInstance", "input.dataset":"someInputDataset", "input.project":"someInputProject"}

start_pipeline = CloudDataFusionStartPipelineOperator(
    location=LOCATION,
    pipeline_name=PIPELINE_NAME,
    instance_name=INSTANCE_NAME,
    runtime_args=RUNTIME_ARGS,
    task_id="start_pipeline",
)

我的问题是,每次触发 DAG 时,Cloud Composer 都会在 myCloudComposerProject 中查找 myDataFusionInstance 而不是 myDataFusionProject,这会产生类似这样的错误:

googleapiclient.errors.HttpError: <HttpError 404 when requesting https://datafusion.googleapis.com/v1beta1/projects/myCloudComposerProject/locations/someLocation/instances/myDataFusionInstance?alt=json returned "Resource 'projects/myCloudComposerProject/locations/someLocation/instances/myDataFusionInstance' was not found". Details: "[{'@type': 'type.googleapis.com/google.rpc.ResourceInfo', 'resourceName': 'projects/myCloudComposerProject/locations/someLocation/instances/myDataFusionInstance'}]"

所以问题是:如何强制我的操作员使用 Data Fusion 项目而不是 Cloud Composer 项目?我怀疑我可以通过添加新的运行时参数来做到这一点,但我不知道该怎么做。

最后一条信息:数据融合管道只是从 BigQuery 源中提取数据并将所有内容发送到 BigTable 接收器。

【问题讨论】:

  • 我认为您可以在使用CloudDataFusionStartPipelineOperator 时指定project_id。在code 上滚动到CloudDataFusionStartPipelineOperator,你会发现你可以设置project-id。你试过了吗?
  • 我最终也注意到了那个参数,我不知道为什么我错过了它......现在它起作用了:) 你能把你的想法写成官方答案,以便我验证它吗?

标签: airflow google-cloud-composer google-cloud-data-fusion


【解决方案1】:

作为在气流上开发运算符时的建议,我们应该检查实现运算符的类,因为文档可能由于版本控制而缺少一些信息。

正如评论的那样,如果您检查CloudDataFusionStartPipelineOperator,您会发现它使用了一个挂钩来获取基于project-id 的实例。此项目 ID 是可选的,因此您可以添加自己的 project-id

class CloudDataFusionStartPipelineOperator(BaseOperator):
 ...

    def __init__(
       ...
        project_id: Optional[str] = None,   ### NOT MENTION IN THE DOCUMENTATION 
        ...
    ) -> None:
        ...
        self.project_id = project_id 
        ...

    def execute(self, context: dict) -> str:
        ...
        instance = hook.get_instance(
            instance_name=self.instance_name,
            location=self.location,
            project_id=self.project_id, ### defaults your project-id
        )
        api_url = instance["apiEndpoint"]
        ... 

将参数添加到您的操作员调用应该可以解决您的问题。

start_pipeline = CloudDataFusionStartPipelineOperator(
    location=LOCATION,
    pipeline_name=PIPELINE_NAME,
    instance_name=INSTANCE_NAME,
    runtime_args=RUNTIME_ARGS,
    project_id=PROJECT_ID,
    task_id="start_pipeline",
)

最后一点,除了official documentation site,您还可以在github 上浏览apache 气流文件。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-05-27
    • 1970-01-01
    • 2021-08-11
    • 1970-01-01
    • 2022-11-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多