标签: google-cloud-composer

Google Cloud Composer和Google Cloud SQL

我们有哪些方法可以从新推出的Google Cloud Composer连接到Google Cloud SQL(MySQL)实例?目的是将Cloud SQL实例中的数据导入BigQuery(可能通过云存储实现中间步骤).

  1. Cloud SQL代理是否可以在托管上以某种方式暴露给托管Composer的Kubernetes集群?

  2. 如果没有,可以使用Kubernetes Service Broker引入Cloud SQL Proxy吗?- > https://cloud.google.com/kubernetes-engine/docs/concepts/add-on/service-broker

  3. 应该使用Airflow来安排和调用GCP API命令,例如1)将mysql表导出到云存储2)读取mysql导出到bigquery?

  4. 也许还有其他方法让我无法完成这项工作

google-cloud-sql google-cloud-platform airflow google-cloud-composer

6
推荐指数
2
解决办法
2556
查看次数

是否可以在 Google Cloud Composer 上安装 github 存储库

如标题所示,我们可以在requirements.txt文件中设置pypi包并使用命令

gcloud beta composer environments update env_name --update-pypi-packages-from-file requirements.txt --location location
Run Code Online (Sandbox Code Playgroud)

更新 Cloud Composer 环境。

但它是否支持在requirements.txt中安装自定义github repo?我尝试添加链接,例如:

pkg_name @ git+ssh://git@github.com/my_account/pkg_repo.git#master
Run Code Online (Sandbox Code Playgroud)

但它不起作用。

谢谢!

更新: 我有一个解决方法是将库放入插件中。但我认为在我们的例子中最好的策略是从 github 安装一个包。

github google-cloud-platform airflow google-cloud-composer

6
推荐指数
1
解决办法
1675
查看次数

使用插件导入DAG时出现气流错误 - 只能在操作员之间设置关系

我编写了一个气流插件,它只包含一个自定义运算符(以支持BigQuery中的CMEK).我可以使用单个任务创建一个简单的DAG,该任务使用此运算符并且执行正常.

但是,如果我尝试在DAG中从DummyOperator任务创建依赖关系到我的自定义操作员任务,则DAG无法在UI中加载并抛出以下错误,我无法理解为什么会抛出此错误?

破坏的DAG:[/home/airflow/gcs/dags/js_bq_custom_plugin_v2.py]关系只能在运营商之间设置; 收到BQCMEKOperator

到目前为止,我已经在composer-1.4.2-airflow-1.9.0,composer-1.4.2-airflow-1.10.0和composer-1.4.1-airflow-1.10.0上进行了测试.

每个任务的运行气流测试都可以顺利完成.

在DAG中单独使用它可以正常工作(如下所示)所以我不相信插件本身存在任何错误

import datetime
import logging
from airflow.models import DAG
from airflow.operators.bq_cmek_plugin import BQCMEKOperator


default_dag_args = {
    'start_date': datetime.datetime(2019,1,1),
    'retries': 0
}


dag = DAG(
    'js_bq_custom_plugin',
    schedule_interval=None,
    catchup=False,
    concurrency=1,
    max_active_runs=1,
    default_args=default_dag_args)

run_this = BQCMEKOperator(
    task_id     = 'cmek_plugin_test',
    sql         = 'select * from ds.foo LIMIT 15',
    project     = 'xxx',
    dataset     = 'js_dev',
    table       = 'cmek_test10',
    key         = 'xxx',
    dag     = dag
)
Run Code Online (Sandbox Code Playgroud)

然而,如果我引入DummyOperator和依赖项,则会发生错误

import datetime
import logging
from airflow.models import DAG
from airflow.operators.bq_cmek_plugin import BQCMEKOperator
from …
Run Code Online (Sandbox Code Playgroud)

airflow google-cloud-composer

6
推荐指数
1
解决办法
1183
查看次数

Airflow:触发 DAG 运行时出现重复条目​​ mysql 完整性错误

我有两个 Airflow DAG - 调度程序和工作人员。调度程序每分钟运行一次,轮询新的聚合作业并触发辅助作业。您可以在下面找到调度程序作业的代码。

然而,在 6000 多个调度程序作业中,有 30 个运行失败,异常情况如下:

[2019-05-14 11:02:12,382] {models.py:1760} ERROR - (MySQLdb._exceptions.IntegrityError) (1062, "Duplicate entry 'run_query-worker-2019-05-14 11:02:11.000000' for key 'PRIMARY'") [SQL: 'INSERT INTO task_instance (task_id, dag_id, execution_date, start_date, end_date, duration, state, try_number, max_tries, hostname, unixname, job_id, pool, queue, priority_weight, operator, queued_dttm, pid, executor_config) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)'] [parameters: ('run_query', 'worker', datetime.datetime(2019, 5, 14, 11, 2, 11, tzinfo=<Timezone [UTC]>), None, None, …
Run Code Online (Sandbox Code Playgroud)

google-cloud-platform airflow airflow-scheduler google-cloud-composer

6
推荐指数
1
解决办法
1944
查看次数

Cloud Composer 上的 Airflow 无法导入模块

我正在test_dag.py我的 Google Cloud Storage Bucket 中运行一个 DAG,其结构如下。

gcs-bucket/
    dags/
        test_dag.py
        dependencies/
            __init__.py
            dependency_1.py
            module1/
                __init__.py
                dependency_2.py
Run Code Online (Sandbox Code Playgroud)

Airflow 检测到 DAG, test_dag.py,它尝试从 导入depencies/dependency_1.py,(导入成功)并dependencies/module1/dependency_2.py给出错误Broken DAG: [/home/airflow/gcs/dags/test_dag.py] module 'dependencies' has no attribute 'module1'

导致此问题的线路是from dependencies.module1 import dependency_2.

这似乎向我表明 Cloud Composer 无法从 中的子目录导入,并且当我在此处dependencies/查看它们的依赖项文档时,他们给出的示例仅是下一级目录(并且只有 1 个文件,而不是完整的文件) python 包)。/dags

不过,这是一个奇怪的部分——当我在 Airflow 本地(而不是在 Cloud Composer 上)运行它时,它运行成功。因此,我不知道为什么我的导入可以在本地运行,但不能在 Cloud Composer 上运行。

我还尝试从我的__init__.py文件中导入所有内容,这给了我相同的属性错误,并将我的依赖项上移到gcs-bucket/似乎根本找不到它们的位置。

当我__file__用 DAG/home/airflow/gcs/dags/test_dag.py打印时sys.path …

python airflow google-cloud-composer

6
推荐指数
1
解决办法
3041
查看次数

Airflow SFTPHook - 找不到主机的主机密钥

我尝试通过传入 a 来使用 Airflow SFTPHookssh_conn_id,但出现错误:

No hostkey for host myhostname found.
Run Code Online (Sandbox Code Playgroud)

然而,使用SFTPOperator进行同样的操作ssh_conn_id却可以正常工作。我该如何解决这个错误?

airflow google-cloud-composer

5
推荐指数
1
解决办法
3775
查看次数

GKE 自动扩缩不会缩小规模

我们使用 GKE(Google Kubernetes Engine)在 GCC(Google Cloude Composer)中运行 Airflow 作为我们的数据管道。

我们一开始有 6 个节点,后来意识到成本飙升,而且我们没有使用那么多的 CPU。所以我们认为我们可以降低最大值,但也可以启用自动缩放。

由于我们在夜间运行管道,并且白天只运行较小的作业,因此我们希望在 1-3 个节点之间运行自动缩放。

因此,我们在 GKE 节点池上启用了自动缩放,但没有按照他们的建议在 GCE 实例组上启用自动缩放。然而,我们得到这个: 节点池无法扩展

为什么是这样?

下面是过去 4 天的 CPU 利用率图表: 在此输入图像描述

我们从未超过 20% 的使用率,那为什么不缩小规模呢?

今天早上我们手动将其缩小到 3 个节点。

google-compute-engine google-cloud-platform google-kubernetes-engine google-cloud-composer

5
推荐指数
1
解决办法
4260
查看次数

如何将单个文件作为卷挂载到 KubernetesPodOperator?

我有一个 docker 映像,需要在启动时安装 JSON 凭证文件。容器通过如下命令启动:

docker run -v [CREDENTIALS_FILE]:/credentials.json image_name

该映像位于 Google 容器注册表中,我想使用 KubernetesPodOperator 在 Cloud Composer dag 中启动它。

有没有办法通过 KubernetesPodOperator 挂载单个文件?理想情况下,该文件将托管在云存储位置。我读到有一个volume/volume_mount选项,但传递单个文件似乎是一件繁重的事情——希望还有另一个我忽略的选项。

KubernetesPodOperator(namespace='default',
                      image="gcr.io/image_name,
                      name="start-container-image",
                      task_id="start-container-image",
                      volume=[?],
                      volume_mounts=[?],
                      dag=dag)
Run Code Online (Sandbox Code Playgroud)

google-kubernetes-engine airflow google-cloud-composer

5
推荐指数
1
解决办法
2698
查看次数

如何在 Airflow 中使用电子邮件操作符附加文件

我使用了参数 files =["abc.txt"]。我从气流文档中获取了信息... https://airflow.readthedocs.io/en/stable/_modules/airflow/operators/email_operator.html

但我收到找不到该文件的错误。我的问题是这个气流将从哪里选择我的文件。是来自 Composer 环境中的 GCS Bucket 还是 DAG 文件夹?

我需要在哪里上传文件以及“文件”参数的正确语法是什么?

提前致谢。

google-cloud-platform airflow google-cloud-composer

5
推荐指数
1
解决办法
2万
查看次数

Airflow:为什么 DAG 任务运行过时的 DAG 代码?

我正在 GCP 上通过 Cloud Composer (1.11.1) 运行 Airflow (1.10.9)。每当我更新 DAG 的代码时,我都可以在 Airflow GUI 中看到更新后的代码已刷新,但至少有 10 分钟 DAG 的任务仍然运行旧代码。

有几个问题:

  1. 为什么会出现这种延迟?可以减少这种延迟吗?
  2. 我如何知道任务的代码何时已更新以确保没有人运行旧代码?

airflow google-cloud-composer

5
推荐指数
0
解决办法
993
查看次数