【问题标题】:mongoimport JSON from Google Cloud Storage in an Airflow task在 Airflow 任务中从 Google Cloud Storage 导入 JSON
【发布时间】: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


    【解决方案1】:

    您无需编写 shell 脚本即可从 GCS 下载。您可以简单地使用GCSToLocalFilesystemOperator,然后您可以使用MongoHook 的insert_many 函数打开文件并将其写入mongo。

    我没有测试它,但它应该是这样的:

    mongo = MongoHook(conn_id=mongo_conn_id)
    with open('file.json') as f:
        file_data = json.load(f)
    mongo.insert_many(file_data)
    

    这是一个管道:BigQuery -> GCS -> 本地文件系统 -> MongoDB。

    如果您愿意,也可以在内存中执行以下操作:BigQuery -> GCS -> MongoDB。

    【讨论】:

    • 这是有道理的。到目前为止,我还没有在我的 Airflow 项目中使用过钩子,但我会尝试一下。我以前使用PythonOperatorpandas,通过将BigQuery 中的数据查询到pandas 数据框中,然后将insert_many'ing 到MongoDB,(我认为)它都在内存中。我希望 (a) 避免内存不足问题并 (b) 提高​​管道的速度,这就是为什么我现在尝试在不使用 pandas 数据帧作为中介的情况下执行 BigQuery -> GCS -> MongoDB。
    • 有一个s3_to_mongo钩子,但是好像没有gcs_to_mongo钩子,这很糟糕。
    • 我也打算用mongoimport代替python的mongo.insert_many函数,希望它能提高插入我们mongo集群的速度
    • @Canovic 我没有在气流回购中看到 s3_to_mongo。也许您正在寻找某人的自定义代码?至于您的问题:(a)避免内存问题的方法是在内存中做更少的事情,这意味着使用 GCS -> 本地文件系统 -> MongoDB。我不知道你的管道,但通常的 ETL 不仅仅是移动东西,它也是关于转换,所以如果你想在 GCS 或本地磁盘上转换,这是你的电话。还有GCSFileTransformOperator。 (b) 加速 ETL 可以来自许多领域。你能并行任务吗?例如,您只复制 1 个文件还是多个文件?考虑动态任务。
    • @Canovic 这不是官方的气流回购。这是属于天文学家的回购。他们分享了一些自定义运算符来帮助社区。我认为您描述的 ELT 相对简单。 “困难部分”只是本地文件系统 -> MongoDB,这是您需要自己编写的东西 - 或者更准确地说,您只需要利用钩子功能。
    猜你喜欢
    • 2020-01-16
    • 2022-08-15
    • 2019-02-14
    • 2023-04-08
    • 2018-04-29
    • 2021-07-27
    • 1970-01-01
    • 2018-06-14
    • 2014-03-24
    相关资源
    最近更新 更多