Memory leak in distributed.worker dask

Viewed 816

I am using distributed worker in dask to make ugly parallel computation.

import dask.dataframe as dd
import json
import os
import dask.bag as db
import pandas as pd
from dask.distributed import Client

def decode_utf8(input_string):
    return json.loads(input_string.encode().decode("utf-8-sig"))

def read_json_file(file_path):
    return db.read_text(file_path).map(decode_utf8)

def batch_process_id(local_path, output_path):
    start_time = time.time()
    client = Client()

    print("-" * 30)
    dask_df = read_json_file(local_path).to_dataframe()
    dask_df.repartition(npartitions=160)
    dask_ecounter_id_schema = ("dummy_id","int")


    temp_series = dask_df["Encounter"].apply(lambda row: int(row["Id"]) if "Id" in row else -1,
                                             meta = dask_ecounter_id_schema)
    temp_series.compute()


if __name__ == "__main__":
    local_path_list =  ["./sessions-201906*.txt",
                       ]
    output_path_list = ["./flowsheet06"]

    for i in range(len(local_path_list)):
        batch_process_id(local_path_list[i], output_path_list[i])

I am using 4-core cluster with total memory of 16 GB. The file size varies between 0.4GB to 0.9GB in sessions-201906*.txt and about 40 files.

When I ran the program, I encountered these issues: 1. istributed.worker - WARNING - Memory use is high but worker has no data to store to disk. Perhaps some other process is leaking memory? Process memory: 3.27 GB -- Worker memory limit: 4.08 GB 2. distributed.utils_perf - WARNING - full garbage collections took 34% CPU time recently (threshold: 10%) distributed.worker - WARNING - gc.collect() took 1.904s. This is usually a sign that some tasks handle too many Python objects at the same time. Rechunking the work into smaller tasks might help.

Can anyone provide any suggestions to handle these issues?

0 Answers
Related