import logging
from time import sleep
from typing import List

from prefect import task, Flow, unmapped, flatten
from prefect.engine.executors import LocalDaskExecutor

@task
def individual_preprocessing(name: str) -> str:
    msg = 'Preprocessing ' + name
    logging.info(msg)
    print(msg)
    sleep(5)
    return name.upper()

@task
def spread_out(data: str, letters: list) -> List[str]:
    sleep(2)
    print("spread_out: {}, {}".format(data, letters))
    return [(data, elem) for elem in letters]

@task
def pront(elem):
    print("elem: {}".format(elem))

def create_workflow(data: list, letters: str):

    with Flow('demo') as flow:
        obj = individual_preprocessing.map(letters)
        spread_mapped = spread_out.map(data, unmapped(obj))
        pront.map(flatten(spread_mapped))

    return flow


def main():
    flow = create_workflow([0, 1, 2, 3], 'abc')

    # flow.visualize()
    flow.run(executor=LocalDaskExecutor(scheduler='threads', n_threads=8))


if __name__ == '__main__':
    logging.basicConfig(level=logging.INFO, format='[%(asctime)s] %(message)s ', filename='demo.log')
    main()