Convert several json files into a single dask dataframe and save this dataframe in a database

Viewed 356

I have a large number of json files in a folder. I would like to read them and save them in a single database. I had thought of pandas dataframes to do this but due to the large number of files, this operation is very slow. Someone suggested dask, which also has dataframes and seems to be very fast, so I installed it and did some tests. It seems to be confirmed. but unlike pandas, the dataframe just contains the content of each json as text. Can someone tell me how to do this so that I get a single dataframe with all the json files read and with the column names of the keys of each object in each json as well as the values?

pandas code

import pandas as pd


data = []
folder = '20-05-2019'
json_files = get_json_files(folder)
for json_file in json_files:
    df_temp = pd.read_json(json_file, encoding='utf-8')
    data.append(df_temp)

df = pd.concat(data)

df.head(10)

result

https://i.imgur.com/NMLH5dj.png

dask dataframe code

code+result

sample folder with 7000 json files in each folder (3000 folder in total)

https://i.imgur.com/JxzYO6X.png

sample json

[
    {
        "sector":"AAA",
        "code":"0009",
        "id":"00000001",
        "fname":"FirstName",
        "lname":"LastName",
        "height":"158",
        "dob":"01/03/2006",
        "din":"19/05/2019 13:23",
        "dout":"19/05/2019 17:46",
        "type":"some text",
        "group":"2",
        "dod":"19/05/2019 13:48",
        "desc":"some text",
        "details":"some text",
        "feval":"Triage 1",
        "localisation":"some place",
        "infop":"not yet implemented",
        "is_in":true,
        "is_op":false,
        "list_c":[],
        "list_e":[],
        "list_a":[],
        "list_r":[]
    },
    {
        "sector":"AAA",
        "code":"0009",
        "id":"00000001",
        "fname":"FirstName",
        "lname":"LastName",
        "height":"158",
        "dob":"01/03/2006",
        "din":"19/05/2019 13:23",
        "dout":"19/05/2019 17:46",
        "type":"some text",
        "group":"2",
        "dod":"19/05/2019 13:48",
        "desc":"some text",
        "details":"some text",
        "feval":"Triage 1",
        "localisation":"some place",
        "infop":"not yet implemented",
        "is_in":true,
        "is_op":false,
        "list_c":[],
        "list_e":[],
        "list_a":[],
        "list_r":[]
    },
    {
        "sector":"AAA",
        "code":"0009",
        "id":"00000001",
        "fname":"FirstName",
        "lname":"LastName",
        "height":"158",
        "dob":"01/03/2006",
        "din":"19/05/2019 13:23",
        "dout":"19/05/2019 17:46",
        "type":"some text",
        "group":"2",
        "dod":"19/05/2019 13:48",
        "desc":"some text",
        "details":"some text",
        "feval":"Triage 1",
        "localisation":"some place",
        "infop":"not yet implemented",
        "is_in":true,
        "is_op":false,
        "list_c":[],
        "list_e":[],
        "list_a":[],
        "list_r":[]
    },
    {
        "sector":"AAA",
        "code":"0009",
        "id":"00000001",
        "fname":"FirstName",
        "lname":"LastName",
        "height":"158",
        "dob":"01/03/2006",
        "din":"19/05/2019 13:23",
        "dout":"19/05/2019 17:46",
        "type":"some text",
        "group":"2",
        "dod":"19/05/2019 13:48",
        "desc":"some text",
        "details":"some text",
        "feval":"Triage 1",
        "localisation":"some place",
        "infop":"not yet implemented",
        "is_in":true,
        "is_op":false,
        "list_c":[],
        "list_e":[],
        "list_a":[],
        "list_r":[]
    },
    {
        "sector":"AAA",
        "code":"0009",
        "id":"00000001",
        "fname":"FirstName",
        "lname":"LastName",
        "height":"158",
        "dob":"01/03/2006",
        "din":"19/05/2019 13:23",
        "dout":"19/05/2019 17:46",
        "type":"some text",
        "group":"2",
        "dod":"19/05/2019 13:48",
        "desc":"some text",
        "details":"some text",
        "feval":"Triage 1",
        "localisation":"some place",
        "infop":"not yet implemented",
        "is_in":true,
        "is_op":false,
        "list_c":[],
        "list_e":[],
        "list_a":[],
        "list_r":[]
    },
    {
        "sector":"AAA",
        "code":"0009",
        "id":"00000001",
        "fname":"FirstName",
        "lname":"LastName",
        "height":"158",
        "dob":"01/03/2006",
        "din":"19/05/2019 13:23",
        "dout":"19/05/2019 17:46",
        "type":"some text",
        "group":"2",
        "dod":"19/05/2019 13:48",
        "desc":"some text",
        "details":"some text",
        "feval":"Triage 1",
        "localisation":"some place",
        "infop":"not yet implemented",
        "is_in":true,
        "is_op":false,
        "list_c":[],
        "list_e":[],
        "list_a":[],
        "list_r":[]
    },
]

EDIT for further information

Maybe I didn't explain well enough why I put everything in a dataframe first. In fact, in each json a certain amount of information is stored about people and this data is updated at regular intervals. At the end of the day, the backup operation in the database takes place. So it is sometimes and often a question of duplicate data between files that must be processed before saving everything in the database.

This means that the files cannot be processed individually. You have to group them all together in a single dataframe (this is the best idea I have at the moment).

1 Answers

As often happens, the answer is actually simple - you do not need any of dask's dataframe API for this, you will be acting on your files one-by-one. This requires no interaction between the tasks, and no output, so you can just use the delayed function.

from dask import compute, delayed

@dask.delayed
def process(file_name):
    df = pd.read_json(json_file, encoding='utf-8')
    df.to_sql(...)

tasks = [process(fn) for fn in json_files]
dask.compute(*tasks)

This probably makes heavy use of the python interpreter, so you will find the above cannot use all the cores when in a single process. I would recommend you use dask.distributed to create a local cluster with as many processes as you have cores, and one thread per process.

Related