【发布时间】:2020-11-14 01:25:34
【问题描述】:
似乎将数据从 GCS 移动到 MongoDB 并不常见,因为这方面的文档并不多。我们将以下任务作为 python_callable 传递给 Python 运算符 - 此任务将数据从 BigQuery 以 JSON 格式移动到 GCS:
def transfer_gcs_to_mongodb(table_name):
# connect
client = bigquery.Client()
bucket_name = "our-gcs-bucket"
project_id = "ourproject"
dataset_id = "ourdataset"
destination_uri = f'gs://{bucket_name}/{table_name}.json'
dataset_ref = bigquery.DatasetReference(project_id, dataset_id)
table_ref = dataset_ref.table(table_name)
configuration = bigquery.job.ExtractJobConfig()
configuration.destination_format = 'NEWLINE_DELIMITED_JSON'
extract_job = client.extract_table(
table_ref,
destination_uri,
job_config=configuration,
location="US",
) # API request
extract_job.result() # Waits for job to complete.
print("Exported {}:{}.{} to {}".format(project_id, dataset_id, table_name, destination_uri))
此任务已成功将数据导入 GCS。但是,当谈到如何正确运行 mongoimport 以将这些数据导入 MongoDB 时,我们现在陷入了困境。特别是mongoimport好像不能指向GCS中的文件,而是要先下载到本地,再导入MongoDB。
这应该如何在 Airflow 中完成?我们是否应该编写一个从 GCS 下载 JSON 的 shell 脚本,然后使用正确的 uri 和所有正确的标志运行 mongoimport?或者还有其他我们缺少的在 Airflow 中运行 mongoimport 的方法吗?
【问题讨论】:
标签: google-cloud-storage airflow mongoimport