删除VM(通过BashOperator)以便在下次触发airflow dag时重复前面的步骤。
我们只需要以下 bash 命令:
bash_cmd = {
'active_account':
'gcloud auth activate-service-account MYCLIENTEMAIL '
'--key-file=/PATH/TO/MY/JSON/SERVICEACCOUNT',
'set_project':
'gcloud config set project MYPROJECTID',
'list_vm':
'gcloud compute instances list',
'create_vm':
'gcloud compute instances create-with-container VMNAME '
'--project=MYPROJECTID --zone=MYZONE --machine-type=e2-medium '
'--image=projects/cos-cloud/global/images/cos-stable-101-17162-40-5 '
'--boot-disk-size=10GB --boot-disk-type=pd-balanced '
'--boot-disk-device-name=VMNAME '
'--container-image=eu.gcr.io/MYPROJECTID/MYCONTAINER --container-restart-policy=always '
'--labels=container-vm=cos-stable-101-17162-40-5 --no-shielded-secure-boot '
'--shielded-vtpm --shielded-integrity-monitoring '
'--metadata startup-script="#!/bin/bash
sleep 10
sudo useradd -m bob
sudo -u bob docker-credential-gcr configure-docker
sudo usermod -aG docker bob
sudo -u bob docker run eu.gcr.io/MYPROJECTID/MYCONTAINER
sudo poweroff" ',
'delete_vm':
'gcloud compute instances delete VMNAME --zone=MYZONE --delete-disks=boot',
}
active_account 和 set_project 分别用于激活服务帐户并设置正确的工作项目(我们要在其中运行 VM)。当 Airflow 在实例化 VM 的 GCP 项目之外运行时,这是必需的。对所使用的服务帐户拥有 ComputeEngine 权限也很重要。要运行的容器映像必须位于实例化 VM 的同一项目的容器注册表中。
list_vm 返回项目中现有 VM 的列表以及相关功能和状态 (RUNNING/STOPPED)。
create_vm 创建附加 docker 以从容器注册表运行的 VM。可以根据需要自定义创建 VM 的命令。需要注意的是,您必须添加 --metadata startup-script ,其中包括 docker 的运行和 docker 完成运行时 VM 的电源关闭。 (要查看如何生成启动脚本,请参阅here)。
delete_vm 只是删除create_vm 创建的虚拟机。
所有这些命令都可以通过以下方式组合在一起在 Airflow DAG 中协同工作:
import re
import os
import datetime
import subprocess
import airflow
from airflow.sensors.python import PythonSensor
from airflow.operators.bash_operator import BashOperator
def vm_run_check():
"function to list all the VMs and check their status"
finish_run = False
output = subprocess.check_output(
bash_cmd['active_account'] + " && " +
bash_cmd['set_project'] + " && " +
bash_cmd['list_vm'],
shell=True
)
output = output.decode("utf-8").split("
")[:-1]
machines = []
for i in range(1,len(output)):
m = {}
for match in re.finditer(r"([A-Z_]+)( +)?", output[0]+" "*10):
span = match.span()
m[match.group().strip()] = output[i][span[0]:span[1]].strip()
machines.append(m)
machines = {m['NAME']:m for m in machines}
if VMNAME in machines:
if machines[VMNAME]['STATUS'] == 'TERMINATED':
finish_run = True
return finish_run
default_args = {
'owner': 'airflow',
'depends_on_past': False,
'email': [''],
'email_on_failure': False,
'email_on_retry': False,
'retries': 0,
}
with models.DAG(
'MYDAGNAME',
catchup=False,
default_args=default_args,
start_date=datetime.datetime.now() - datetime.timedelta(days=3),
schedule_interval='0 4 * * *', # every day at 04:00 AM UTC
) as dag:
create_vm = BashOperator(
task_id="create_vm",
bash_command = bash_cmd['active_account'] + " && " +
bash_cmd['set_project'] + " && " +
bash_cmd['create_vm']
)
sensor_vm_run = PythonSensor(
task_id="sensor_vm_run"
python_callable=vm_run_check,
poke_interval=60*2, # check every 2 minutes
timeout=60*60, # check every 2 minutes for an hour
soft_fail=True,
mode="reschedule",
)
delete_vm = BashOperator(
task_id="delete_vm",
bash_command = bash_cmd['active_account'] + " && " +
bash_cmd['set_project'] + " && " +
bash_cmd['delete_vm']
)
create_vm >> sensor_vm_run >> delete_vm