import dask
from dask.distributed import Client
import prefect
from prefect.engine.executors import LocalDaskExecutor, DaskExecutor
from sklearn.linear_model import LinearRegression
import dask_ml
from dask_ml.datasets import make_classification, make_regression
import pandas as pd
import numpy as np

client = Client()

client

X, y = make_regression(n_samples=500_000, n_features=100, random_state=42, chunks=100_000)
X

pd.DataFrame(X.compute()).to_csv("X_test", header=False, index=False, sep="\t")
pd.DataFrame(y.compute()).to_csv("y_test", header=False, index = False, sep="\t")

np.random.seed(42)

@prefect.task
def read_data(file1:str, file2:str):
    X = pd.read_csv("X_test", header=None, sep="\t")
    y = pd.read_csv("y_test", header=None, sep="\t")
    return X, y

@prefect.task
def make_random_iter(num_batches:int):
    return np.random.rand(num_batches)

@prefect.task
def fit_create_model(X_y, rand:float):
    X,y = X_y
   
    coef_list =[]
    model = LinearRegression()
    for i in range(10):
        y = y + np.random.rand()
        coef_list.append(model.fit(X,y).coef_[0])
    return coef_list

with prefect.Flow("test") as flow:
    X_file = prefect.Parameter("X_file")
    y_file = prefect.Parameter("y_file")
    num_100_batches = prefect.Parameter("num_100_batches")
    
    X_y = read_data(X_file, y_file)
    
    rand_list = make_random_iter(num_100_batches)
    coef_list = fit_create_model.map(prefect.unmapped(X_y), rand_list)

run_sequential = flow.run(num_100_batches=10, y_file="y_test", X_file="X_test") #runs fine

executor = DaskExecutor("tcp://127.0.0.1:63558")
run_dask = flow.run(num_100_batches=10, y_file="y_test", X_file="X_test", executor=executor) #fails here