【问题标题】:Python read .json files from GCS into pandas DF in parallelPython 将 GCS 中的 .json 文件并行读取到 pandas DF 中
【发布时间】:2020-07-23 01:02:37
【问题描述】:

TL;DR: asyncio vs multi-processing vs threading vs. some other solution 以并行化从 GCS 读取文件的 for 循环,然后将此数据一起附加到 pandas 数据帧中,然后写入 BigQuery...

我想让一个 python 函数并行化,它从 GCS 目录中读取数十万个小 .json 文件,然后将这些 .jsons 转换为 pandas 数据帧,然后将 pandas 数据帧写入 BigQuery 表。

这是该函数的非并行版本:

import gcsfs
import pandas as pd
from my.helpers import get_gcs_file_list
def load_gcs_to_bq(gcs_directory, bq_table):

    # my own function to get list of filenames from GCS directory
    files = get_gcs_file_list(directory=gcs_directory) # 

    # Create new table
    output_df = pd.DataFrame()
    fs = gcsfs.GCSFileSystem() # Google Cloud Storage (GCS) File System (FS)
    counter = 0
    for file in files:

        # read files from GCS
        with fs.open(file, 'r') as f:
            gcs_data = json.loads(f.read())
            data = [gcs_data] if isinstance(gcs_data, dict) else gcs_data
            this_df = pd.DataFrame(data)
            output_df = output_df.append(this_df)

        # Write to BigQuery for every 5K rows of data
        counter += 1
        if (counter % 5000 == 0):
            pd.DataFrame.to_gbq(output_df, bq_table, project_id=my_id, if_exists='append')
            output_df = pd.DataFrame() # and reset the dataframe


    # Write remaining rows to BigQuery
    pd.DataFrame.to_gbq(output_df, bq_table, project_id=my_id, if_exists='append')

这个函数很简单:

  • 获取['gcs_dir/file1.json', 'gcs_dir/file2.json', ...],GCS中的文件名列表
  • 遍历每个文件名,并且:
    • 从 GCS 读取文件
    • 将数据转换为 pandas DF
    • 附加到一个主要的 pandas DF
    • 每 5K 循环,写入 BigQuery(因为随着 DF 变大,追加会变慢)

我必须在几个 GCS 目录上运行这个函数,每个目录都有大约 500K 文件。由于读取/写入这么多小文件的瓶颈,对于单个目录,此过程将需要约 24 小时...如果我可以使其更加并行以加快速度,那就太好了,因为这似乎是一项任务适合并行化。

编辑:下面的解决方案很有帮助,但我对在 python 脚本中并行运行特别感兴趣。 Pandas 正在处理一些数据清理,使用bq load 会抛出错误。 asynciogcloud-aio-storage 似乎都可能对这项任务有用,可能是比线程或多处理更好的选择...

【问题讨论】:

  • 为什么要这样做?您可以直接使用bq 命令给出GCS 文件夹的路径和bigquery 中的表名。这样会更快
  • 你指的bq命令是什么?
  • 我已经给出了答案,以便遇到同样问题的其他人可以查看它

标签: python pandas parallel-processing google-cloud-storage python-asyncio


【解决方案1】:

与其将并行处理添加到您的 Python 代码中,不如考虑并行多次调用您的 Python 程序。这是一个技巧,它更容易适用于在命令行上获取文件列表的程序。因此,为了这篇文章,让我们考虑更改程序中的一行:

您的线路:

# my own function to get list of filenames from GCS directory
files = get_gcs_file_list(directory=gcs_directory) # 

换行:

files = sys.argv[1:]  # ok, import sys, too

现在,您可以通过这种方式调用您的程序:

PROCESSES=100
get_gcs_file_list.py | xargs -P $PROCESSES your_program

xargs 现在将获取get_gcs_file_list.py 输出的文件名并并行调用your_program 最多100 次,在每行上尽可能多地匹配文件名。我相信文件名的数量仅限于 shell 允许的最大命令大小。如果 100 个进程不足以处理您的所有文件,xargs 将再次调用your_program,直到它从标准输入读取的所有文件名都被处理。 xargs 确保同时运行的 your_program 调用不超过 100 个。您可以根据主机可用的资源来改变进程的数量。

【讨论】:

  • 好主意。感谢分享。我目前在每天运行的 Airflow DAG 中以 tasks 的身份调用我的程序,我不太确定如何将这种模式引入 Airflow。
【解决方案2】:

你可以直接使用bq命令来代替这个。

bq 命令行工具是基于 Python 的 BigQuery 命令行工具。

当您使用此命令时,加载发生在 google 的网络中,这比我们创建数据框并加载到表中要快。

    bq load \
    --autodetect \
    --source_format=NEWLINE_DELIMITED_JSON \
    mydataset.mytable \
    gs://mybucket/my_json_folder/*.json

欲了解更多信息 - https://cloud.google.com/bigquery/docs/loading-data-cloud-storage-json#loading_json_data_into_a_new_table

【讨论】:

  • 我会试试这个。我担心我的数据存在type 问题,尽管您的共享命令似乎有一个--autodetect 标志来处理这个问题?
  • --autodetect 会自动检测类型并尝试应用类型,否则会抛出错误
猜你喜欢
  • 2021-02-02
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-10-14
  • 1970-01-01
相关资源
最近更新 更多