【发布时间】:2019-12-19 04:30:15
【问题描述】:
Airflow 的新手。我正在尝试将结果保存到另一个存储桶(不是气流存储桶)中的文件中。 我可以保存到“/home/airflow/gcs/data/test.json”中的文件,然后使用 gcs_hook.GoogleCloudStorageHook 复制到另一个存储桶。代码如下:
def write_file_func(**context):
file = f'/home/airflow/gcs/data/test.json'
with open(file, 'w') as f:
f.write(json.dumps('{"name":"aaa", "age":"10"}'))
def upload_file_func(**context):
conn = gcs_hook.GoogleCloudStorageHook()
source_bucket = 'source_bucket'
source_object = 'data/test.json'
target_bucket = 'target_bucket'
target_object = 'test.json'
conn.copy(source_bucket, source_object, target_bucket, target_object)
conn.delete(source_bucket, source_object)
我的问题是:
我们可以直接写入目标存储桶的文件吗?我在 gcs_hook 中没有找到任何方法。
我尝试使用google.cloud.storage bucket.blob('test.json').upload_from_string(),但是气流一直说“服务器的DAGBag中没有DAG”,很烦,我们不允许在 DAG 中使用该 API 吗?
如果我们可以直接使用 google.cloud.storage/bigquery API,那和 Airflow API 有什么区别,比如 gcs_hook/bigquery_hook?
谢谢
【问题讨论】:
标签: airflow google-cloud-composer