Tengo un problema de asignación y quería preguntarle a la comunidad SO cuál es la mejor manera de implementar esto para mi marco de datos Spark (utilizando Spark 3.1+). Primero describiré el problema y luego pasaré a la implementación.
Aquí está el problema: tengo hasta N tareas y hasta N individuos (en el caso de este problema, N=10). Cada individuo tiene un costo de realizar cada tarea, donde el costo mínimo es $0 y el costo máximo es $10. Es una especie de problema de algoritmo húngaro con algunas advertencias.
multiTask=True (no puede haber más de 1 multiTask y es posible que no haya ninguno). Si un trabajador tiene un costo menor que x para la multitarea, se le asigna automáticamente la multitarea y la multitarea se considera tomada durante la optimización.Así es como se ve el marco de datos de Spark. Nota : muestro un ejemplo donde N = 3 (3 tareas, 3 personas) por simplicidad.
from pyspark.sql import Row rdd = spark.sparkContext.parallelize([ Row(date='2019-08-01', locationId='z2-NY', workerId=129, taskId=220, cost=1.50, isMultiTask=False), Row(date='2019-08-01', locationId='z2-NY', workerId=129, taskId=110, cost=2.90, isMultiTask=True), Row(date='2019-08-01', locationId='z2-NY', workerId=129, taskId=190, cost=0.80, isMultiTask=False), Row(date='2019-08-01', locationId='z2-NY', workerId=990, taskId=220, cost=1.80, isMultiTask=False), Row(date='2019-08-01', locationId='z2-NY', workerId=990, taskId=110, cost=0.90, isMultiTask=True), Row(date='2019-08-01', locationId='z2-NY', workerId=990, taskId=190, cost=9.99, isMultiTask=False), Row(date='2019-08-01', locationId='z2-NY', workerId=433, taskId=220, cost=1.20, isMultiTask=False), Row(date='2019-08-01', locationId='z2-NY', workerId=433, taskId=110, cost=0.25, isMultiTask=True), Row(date='2019-08-01', locationId='z2-NY', workerId=433, taskId=190, cost=4.99, isMultiTask=False) ]) df = spark.createDataFrame(rdd) Verá que hay una fecha/ubicación, ya que necesito resolver este problema de asignación para cada agrupación de fecha/ubicación. Estaba planeando resolver esto asignando a cada trabajador y tarea un "índice" basado en sus ID usando dense_rank() y luego usando pandas UDF, completando la matriz N x N numpy basada en los índices e invocando la función linear_sum_assignment . Sin embargo, no creo que este plan funcione debido al caso del segundo borde que presenté con la multitarea.
worker_order_window = Window.partitionBy("date", "locationId").orderBy("workerId") task_order_window = Window.partitionBy("date", "locationId").orderBy("taskId") # get the dense_rank because will use this to assign a worker ID an index for the np array for linear_sum_assignment # dense_rank - 1 as arrays are 0 indexed df = df.withColumn("worker_idx", dense_rank().over(worker_order_window) - 1) df = df.withColumn("task_idx", dense_rank().over(task_order_window) - 1) def linear_assignment_udf(pandas_df: pd.DataFrame) -> pd.DataFrame: df_dict = pandas_df.to_dict('records') # in case there are less than N rows/columns N = max(pandas_df.shape[0], pandas_df.shape[1]) arr = np.zeros((N,N)) for row in df_dict: # worker_idx will be the row number, task idx will be the col number worker_idx = row.get('worker_idx') task_idx = row.get('task_idx') arr[worker_idx][task_idx] = row.get('cost') rids, cids = linear_sum_assignment(n) return_list = [] # now want to return a dataframe that says which task_idx a worker has for r, c in zip(rids, cids): for d in df_dict: if d.get('worker_idx') == r: d['task_assignment'] = c return_list.append(d) return pd.DataFrame(return_list) schema = StructType.fromJson(df.schema.jsonValue()).add('task_assignment', 'integer') df = df.groupBy("date", "locationId").applyInPandas(linear_assignment_udf, schema) df = df.withColumn("isAssigned", when(col("task_assignment") == col("task_idx"), True).otherwise(False))Como puede ver, este caso no cubre la multitarea en absoluto. Me gustaría resolver esto de la manera más eficiente posible para no estar atado a pandas udf o scipy.
No sé nada sobre las bibliotecas que estás usando, así que no puedo ayudarte con el código, pero creo que deberías hacerlo en dos pasos:
El algoritmo húngaro básico solo funciona en matrices de costo cuadradas, y parece que lo ha manejado correctamente rellenando su matriz de costo con 0, pero hay modificaciones del algoritmo que funcionan con matrices rectangulares. Es posible que desee ver si tiene acceso a una de esas alternativas, ya que podría ser significativamente más rápido.