【问题标题】:group the related values in one group将相关值分组为一组
【发布时间】:2021-12-17 06:52:02
【问题描述】:

尝试根据相关记录对列值进行分组

partColumns = (["partnumber","colVal1","colVal2", "colVal3","colVal4","colVal5"])

partrelations = ([("part0","part1","", "","",""),
                  ("part1","","part2", "","part4",""),
                  ("part2","part3", "", "part5","part6","part7"),
                  ("part10","part11","", "","",""),
                  ("part11","part13","part21", "","",""),
                  ("part13","part21","part18", "","part20",""),
                 ])
df_part_groups = spark.createDataFrame(data=partrelations, schema = partColumns) 

试图得到如下输出 -

edges = (df_part_groups
         .withColumnRenamed("partnumber", "src")
         .withColumnRenamed("colVal1", "dst")
        )

vertices = (edges.select("src").distinct()
            .union(edges.select("dst").distinct())
            .withColumnRenamed("src", "id"))
         
#create a graph and find all connected components
g = G.GraphFrame(vertices, edges)
cc = g.connectedComponents()

display(df_part_groups
        .join(cc.distinct(), df_part_groups.device == cc.id)
        .orderBy("component", "partnumber", "colVal1"))

以上是我要整理的内容

谢谢你的帮助!!

【问题讨论】:

  • 您能否详细说明分组算法?这里是感兴趣的行数/列数 - 最简单的分组是i < 3。
  • 分组基于关系。就像part0与part1有关,part1与part2和part4有关。所以 part0 也间接与 part2 和 part4 相关。同样,part2 也与 part3、part5、part6、part7 相关。这就是为什么我们可以将part0、part1、part2、part3、part4、part5、part6、part7组合在一个组G1中。不,行数没有限制,但可以说我们在这个集合中不会有超过 6 列。
  • 好的。你的列标题/标题和空行呢?
  • 列标题无关紧要,对于这个问题。我们可以将它们视为 - 零件编号 | groupnumber,我们可以跳过空行或忽略它们

标签: python pyspark spark-graphx graphframes


【解决方案1】:

我们可以使用集合交集做一个简单的检查来解决这个问题。 (不知道 GraphFrames :()

第 1 步:将所有部分组合到每个行的单个数组中

from pyspark.sql import functions as F
    
df_part_groups1= df_part_groups.withColumn('parts', F.array('partnumber', 'colVal1', 'colVal2', 'colVal3', 'colVal4', 'colVal5')  )

第 2 步:获取 all_parts,它是组合部分列表的列表,因为需要在各个行中确定组。

def clean_lists(plists):
  return [ list(filter(None, pl)) for pl in plists]

all_parts = clean_lists((df_part_groups1.groupBy(F.lit(1)).agg(F.collect_list('parts').alias('parts')).collect())[0].parts)

第 3 步:使用收集的 all_parts 获取组数据

def part_of_existing_group(gps, pl):
  for key in gps.keys():
    if set(gps[key]) & set(pl):
      gps[key] = list(set(gps[key] + pl))      
      return True
      return False   
      
def findGroups(plists): 
  groups = {}    
  index = 1
  for pl in plists:
    if len(groups.keys()) == 0 or (not part_of_existing_group(groups, pl)):
      groups[f'G{index}'] = pl
      index +=1
  return groups  

第 4 步:根据您创建的组图分配组。

 groups = findGroups(all_parts)
    
    @udf
def get_group_val(part):
  for key in groups.keys():
    if part in groups[key]:
      return key
  return -1

df_part_groups2 = df_part_groups1.withColumn('part', F.explode('parts')).dropDuplicates(['part']).where(~F.col('part').like('')).select('part', 'parts').withColumn('Group', get_group_val('part'))

    df_part_groups2.show()
+------+--------------------+-----+
|  part|               parts|Group|
+------+--------------------+-----+
| part0|[part0, part1, , ...|   G1|
| part1|[part0, part1, , ...|   G1|
|part10|[part10, part11, ...|   G2|
|part11|[part10, part11, ...|   G2|
|part13|[part11, part13, ...|   G2|
|part18|[part13, part21, ...|   G2|
| part2|[part1, , part2, ...|   G1|
|part20|[part13, part21, ...|   G2|
|part21|[part11, part13, ...|   G2|
| part3|[part2, part3, , ...|   G1|
| part4|[part1, , part2, ...|   G1|
| part5|[part2, part3, , ...|   G1|
| part6|[part2, part3, , ...|   G1|
| part7|[part2, part3, , ...|   G1|
+------+--------------------+-----+

【讨论】:

  • 似乎是正确的答案,让我周末测试一下,我会相应地更新你。再次感谢。感谢您的帮助
  • 感谢@greenie 的工作。我会看看是否可以将其转换为葡萄框架解决方案。
  • 不客气!如果您也可以发布 GraphFrame 解决方案,那就太好了。
猜你喜欢
  • 2021-09-12
  • 2015-11-16
  • 1970-01-01
  • 2014-06-03
  • 1970-01-01
  • 1970-01-01
  • 2017-12-23
  • 1970-01-01
相关资源
最近更新 更多