我们有哪些方法可以从新推出的Google Cloud Composer连接到Google Cloud SQL(MySQL)实例?目的是将Cloud SQL实例中的数据导入BigQuery(可能通过云存储实现中间步骤).
Cloud SQL代理是否可以在托管上以某种方式暴露给托管Composer的Kubernetes集群?
如果没有,可以使用Kubernetes Service Broker引入Cloud SQL Proxy吗?- > https://cloud.google.com/kubernetes-engine/docs/concepts/add-on/service-broker
应该使用Airflow来安排和调用GCP API命令,例如1)将mysql表导出到云存储2)读取mysql导出到bigquery?
也许还有其他方法让我无法完成这项工作
google-cloud-sql google-cloud-platform airflow google-cloud-composer
如标题所示,我们可以在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 安装一个包。
我编写了一个气流插件,它只包含一个自定义运算符(以支持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 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
我正在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 …
我尝试通过传入 a 来使用 Airflow SFTPHookssh_conn_id,但出现错误:
No hostkey for host myhostname found.
Run Code Online (Sandbox Code Playgroud)
然而,使用SFTPOperator进行同样的操作ssh_conn_id却可以正常工作。我该如何解决这个错误?
我们使用 GKE(Google Kubernetes Engine)在 GCC(Google Cloude Composer)中运行 Airflow 作为我们的数据管道。
我们一开始有 6 个节点,后来意识到成本飙升,而且我们没有使用那么多的 CPU。所以我们认为我们可以降低最大值,但也可以启用自动缩放。
由于我们在夜间运行管道,并且白天只运行较小的作业,因此我们希望在 1-3 个节点之间运行自动缩放。
因此,我们在 GKE 节点池上启用了自动缩放,但没有按照他们的建议在 GCE 实例组上启用自动缩放。然而,我们得到这个:

为什么是这样?
我们从未超过 20% 的使用率,那为什么不缩小规模呢?
今天早上我们手动将其缩小到 3 个节点。
google-compute-engine google-cloud-platform google-kubernetes-engine google-cloud-composer
我有一个 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) 我使用了参数 files =["abc.txt"]。我从气流文档中获取了信息... https://airflow.readthedocs.io/en/stable/_modules/airflow/operators/email_operator.html
但我收到找不到该文件的错误。我的问题是这个气流将从哪里选择我的文件。是来自 Composer 环境中的 GCS Bucket 还是 DAG 文件夹?
我需要在哪里上传文件以及“文件”参数的正确语法是什么?
提前致谢。
我正在 GCP 上通过 Cloud Composer (1.11.1) 运行 Airflow (1.10.9)。每当我更新 DAG 的代码时,我都可以在 Airflow GUI 中看到更新后的代码已刷新,但至少有 10 分钟 DAG 的任务仍然运行旧代码。
有几个问题: