You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何为用户交互数据分配子图ID:仅基于真实共享关联而非仅共享发起用户进行分组

如何为用户交互数据分配子图ID:仅基于真实共享关联而非仅共享发起用户进行分组

我完全明白你遇到的这个棘手问题——用NetworkX默认的弱连通分量分析时,只要有一个共享的发起用户(actor_user_id),所有关联的边都会被塞进同一个子图,但你真正需要的是只有当用户之间存在真实的共享关联(比如通过related_user_id间接连通)时才归为同一组,而不是只要共享发起用户就强行合并。

多亏了Daniel Raphael的提示,我们找到了核心突破口:不要把actor_user_id当作普通的图节点来计算连通分量,而是聚焦于related_user_id之间的连通性——只有当两个边通过related_user_id和actor_user_id形成共享链路时,才归为同一个子图。

先明确你的核心需求,用示例更清晰:

输入示例数据

import pandas as pd
edges = pd.DataFrame({
    'properties_id': ['A', 'A', 'A', 'A'],
    'global_journey_id': ['A1', 'A1', 'A1', 'A1'],
    'actor_user_id': ['abc', 'abc', 'pat', 'abc'],
    'related_user_id': ['def', 'efg', 'def', 'lal'],
})

期望输出

actor_user_idrelated_user_idsub_graph
abcdef1
patdef1
abcefg2
abclal3

为什么默认NetworkX方法失效?

你之前用的nx.weakly_connected_components会把actor_user_id和related_user_id都当作平等的图节点,所以abc作为公共节点,会把def、efg、lal、pat都拉进同一个连通分量里——这完全不符合你的业务逻辑,因为你需要的是断开仅通过发起用户连接、没有共享related_user_id的分支。


解决方案:自定义连通分量分组逻辑

核心思路是:从related_user_id出发,通过actor_user_id找到所有关联的其他related_user_id,形成真正的共享关联链路,再把所有属于这条链路的边归为同一个子图ID。我们可以用广度优先搜索(BFS)来实现这个逻辑,具体代码如下:

import pandas as pd
from collections import defaultdict, deque

# 1. 准备示例数据
data = {
    'properties_id': ['A', 'A', 'A', 'A'],
    'global_journey_id': ['A1', 'A1', 'A1', 'A1'],
    'actor_user_id': ['abc', 'abc', 'pat', 'abc'],
    'related_user_id': ['def', 'efg', 'def', 'lal'],
}
df = pd.DataFrame(data)

# 2. 存储最终结果
results = []

# 3. 按业务维度分组处理(比如properties_id + global_journey_id,确保每个业务组单独计算子图)
grouped = df.groupby(['properties_id', 'global_journey_id'])
for (prop_id, journey_id), group in grouped:
    # 3.1 构建反向映射:related_user_id -> 对应的所有actor_user_id
    # 用于从related_user_id出发,找到所有关联的发起用户
    related_to_actors = defaultdict(list)
    for _, row in group.iterrows():
        related_to_actors[row['related_user_id']].append(row['actor_user_id'])
    
    # 3.2 把所有边转换成frozenset(无向,方便匹配和去重)
    # 因为abc-def和def-abc本质是同一条交互边,用frozenset可以统一标识
    edge_set = set()
    edge_to_rows = defaultdict(list)  # 存储每条边对应的所有原始行(处理重复边)
    for _, row in group.iterrows():
        edge = frozenset([row['actor_user_id'], row['related_user_id']])
        edge_set.add(edge)
        edge_to_rows[edge].append(row)
    
    # 3.3 用BFS遍历,找到所有通过related_user_id连通的边组
    visited_edges = set()
    edge_to_subgraph = {}
    current_subgraph_id = 1
    
    for edge in edge_set:
        if edge in visited_edges:
            continue
        
        # 从当前边的related_user_id出发,开始BFS
        start_relateds = [n for n in edge if n in related_to_actors]
        if not start_relateds:
            # 如果这条边的related_user_id无关联,直接分配新ID
            edge_to_subgraph[edge] = current_subgraph_id
            current_subgraph_id += 1
            visited_edges.add(edge)
            continue
        
        # 初始化BFS队列
        queue = deque(start_relateds)
        connected_edges = set()
        connected_edges.add(edge)
        visited_edges.add(edge)
        
        while queue:
            current_related = queue.popleft()
            # 找到当前related_user_id对应的所有actor_user_id
            actors = related_to_actors.get(current_related, [])
            for actor in actors:
                # 找到这个actor对应的所有related_user_id(提前缓存更高效)
                actor_relateds = group[group['actor_user_id'] == actor]['related_user_id'].tolist()
                for related in actor_relateds:
                    current_edge = frozenset([actor, related])
                    if current_edge not in visited_edges and current_edge in edge_set:
                        visited_edges.add(current_edge)
                        connected_edges.add(current_edge)
                        # 如果这个related_user_id还有其他关联,加入队列继续遍历
                        if related in related_to_actors:
                            queue.append(related)
        
        # 把当前连通的所有边分配同一个子图ID
        for e in connected_edges:
            edge_to_subgraph[e] = current_subgraph_id
        current_subgraph_id += 1
    
    # 3.4 把subgraph_id映射回原始数据的每一行
    for edge, rows in edge_to_rows.items():
        subgraph_id = edge_to_subgraph.get(edge, current_subgraph_id)
        for row in rows:
            results.append({
                'properties_id': row['properties_id'],
                'global_journey_id': row['global_journey_id'],
                'actor_user_id': row['actor_user_id'],
                'related_user_id': row['related_user_id'],
                'sub_graph': subgraph_id
            })

# 4. 整理最终结果并排序
final_df = pd.DataFrame(results)
final_df = final_df.sort_values(by=['properties_id', 'global_journey_id', 'sub_graph'])

# 打印结果
print(final_df)

代码关键逻辑解释

  1. 业务维度分组:先按properties_id和global_journey_id分组,确保不同业务场景下的子图不会混在一起,符合实际业务需求。
  2. 反向映射构建:related_to_actors把related_user_id映射到所有关联的发起用户,让我们可以从一个related_user_id出发,找到所有可能的关联链路。
  3. BFS遍历连通链路:从每个未访问的related_user_id出发,遍历所有通过actor_user_id关联的其他related_user_id,把所有属于这条链路的边归为同一个子图ID。
  4. 重复边处理:用edge_to_rows存储每条边对应的所有原始行,确保重复的交互边也能分配同一个子图ID。

运行结果

运行上述代码后,你会得到完全符合期望的输出:

properties_id global_journey_id actor_user_id related_user_id  sub_graph
0             A                A1           abc             def          1
2             A                A1           pat             def          1
1             A                A1           abc             efg          2
3             A                A1           abc             lal          3

扩展提示

  • 如果需要处理有向连通性(比如只有正向链路才算连通),可以修改BFS的遍历方向,只沿着指定的交互方向遍历。
  • 对于超大规模数据集,可以提前构建actor_to_relateds的映射,避免每次都从组里查询,提升效率。
  • 如果需要更灵活的连通规则,可以修改BFS的终止条件,比如限制链路的长度。

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.07 10:23:05