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

千万级RDF多表哈希连接报Out of Memory错误优化问询

1000万行RDF三元组多表哈希连接内存溢出优化

问题背景

我有一份存储为.txt格式的数据集,共包含1000万行RDF三元组结构记录,数据样例如下:

wsdbm:User0 wsdbm:follows   wsdbm:User300 .
wsdbm:User6 wsdbm:likes wsdbm:Product92 .
wsdbm:Product0  rev:hasReview   wsdbm:Review478 .
wsdbm:User2 wsdbm:friendOf  wsdbm:User119 .
....

本次使用的RDF三元组结构组成如下:

Subjects: User0, User6, Product, User2
Predicates: follows, likes, hasReview, friendOf
Objects: User300, Product92, Review478, User119

需要执行的查询逻辑对应SQL如下:

SELECT follows.subject, follows.object, friendOf.object, 
       likes.object, hasReview.object
FROM follows, friendOf, likes, hasReview
WHERE follows.object = friendOf.subject
      AND friendOf.object = likes.subject
      AND likes.object = hasReview.subject

现有实现

目前已实现PropertyTables类,内置方法可遍历原始文件,将每个subject、predicate、object转换为整数类型,用于提升连接计算效率、降低内存占用,相关代码如下:

from collections import defaultdict

class PropertyTables():
    """
    存储查询所需的4张属性表,每张表为PropertyTable类实例
    """
    def __init__(self):
        self.property_tables = defaultdict()
        self.hash_map = HashDict()

    def parse_file(self, file_path, remove_prefix = False):
        data = open(file_path, 'r')
        for line in data:
            subj, prop, *obj = line.rstrip('\n.').split('\t')
            obj = obj[0].rstrip()
    
            if remove_prefix:
                subj, prop, obj = [self.remove_prefix(s) for s in (subj, prop, obj)]
            
            if prop in ['follows', 'friendOf', 'likes', 'hasReview']:
                self.hash_and_store(subj, prop, obj)
    
        data.close()

上述类引用的PropertyTable类实现如下:

class PropertyTable():
    """
    单张属性表结构,存储所有对应谓词的主语、宾语值
    """
    def __init__(self):
        self.table = []

    def insert(self, r, s):
        # 若输入已是元组则直接拼接后存储,否则转为二元组存储,主要用于读入数据阶段初始化
        if type(r) == tuple:
            self.table.append(r + s)
        else:
            self.table.append((r, s))

其中HashDict()为简易字典类,用于对字符串值做整数哈希映射,连接完成后可反向取回原始字符串值。

实现的基础哈希连接算法代码如下:

def hash_join(self, property_1: PropertyTable, index_0, property_2: PropertyTable, index_1):
    ht = defaultdict(list)

    # 为第一张表构建连接键哈希表
    for s in property_1.table:
        ht[s[index_0]].append(s)

    # 遍历第二张表完成匹配连接
    joined_table = PropertyTable()
    for r in property_2.table:
        for s in ht[r[index_1]]:
            joined_table.insert(s, r)

    return joined_table

按照查询连接条件,按顺序调用函数完成多表连接的逻辑如下:

# 连接条件:
# follows.object = friendOf.subject
# friendOf.object = likes.subject
# likes.object = hasReview.subject

join_follows_friendOf = hash_join(pt.property_tables['follows'], 1, pt.property_tables['friendOf'], 0)
join_friendOf_likes = hash_join(join_follows_friendOf, 3, pt.property_tables['likes'], 0)
join_likes_hasReview = hash_join(join_friendOf_likes, 5, pt.property_tables['hasReview'], 0)

现存问题

上述逻辑在小规模数据集上运行结果正确,但处理1000万行规模数据时会触发内存溢出(Out of Memory)错误。小数据量下的内存profiling结果如下:

Line #    Mem usage    Increment  Occurrences   Line Contents
=============================================================
    53     68.0 MiB     68.0 MiB           1       @profile
    54                                             def hash_and_store(self, subj, prop, obj):
    55                                         
    56     68.0 MiB      0.0 MiB           1           hashed_subj, hashed_obj = self.hash_map.hash_values(subj, obj)
    57                                         
    58     68.0 MiB      0.0 MiB           1           if prop not in self.property_tables: 
    59                                                     self.property_tables[prop] = PropertyTable()
    60     68.0 MiB      0.0 MiB           1           self.property_tables[prop].insert(hashed_subj, hashed_obj)


Line #    Mem usage    Increment  Occurrences   Line Contents
=============================================================
    32     68.1 MiB     68.1 MiB           1       @profile
    33                                             def parse_file(self, file_path, remove_prefix = False):
    34                                         
    35     68.1 MiB      0.0 MiB           1           data = open(file_path, 'r')
    36                                         
    37                                                 
    38                                                 
    39                                                
    40                                         
    41     80.7 MiB      0.3 MiB      109311           for line in data:
    42     80.7 MiB      0.0 MiB      109310               subj, prop, *obj = line.rstrip('\n.').split('\t')
    43     80.7 MiB      0.5 MiB      109310               obj = obj[0].rstrip()
    44                                         
    45     80.7 MiB      0.0 MiB      109310               if remove_prefix:
    46     80.7 MiB      9.0 MiB      655860                   subj, prop, obj = [self.remove_prefix(s) for s in (subj, prop, obj)]
    47                                                     
    48     80.7 MiB      0.0 MiB      109310               if prop in ['follows', 'friendOf', 'likes', 'hasReview']:
    49     80.7 MiB      2.8 MiB       80084                   self.hash_and_store(subj, prop, obj)
    50                                         
    51     80.7 MiB      0.0 MiB           1           data.close()
    
    
   Line #    Mem usage    Increment  Occurrences   Line Contents
=============================================================
    38     80.7 MiB     80.7 MiB           1       @profile
    39                                             def hash_join(self, property_1: PropertyTable, index_0, property_2: PropertyTable, index_1):
    40                                         
    41     80.7 MiB      0.0 MiB           1           ht = defaultdict(list)
    42                                         
    43                                                 # 为第一张表构建哈希表
    44                                         
    45     81.2 MiB      0.0 MiB       31888           for s in property_1.table:
    46     81.2 MiB      0.5 MiB       31887               ht[s[index_0]].append(s)
    47                                         
    48                                                 # 执行连接
    49                                         
    50     81.2 MiB      0.0 MiB           1           joined_table = PropertyTable()
    51                                         
    52    203.8 MiB      0.0 MiB       45713           for r in property_2.table:
    53    203.8 MiB      0.0 MiB     1453580               for s in ht[r[index_1]]:
    54    203.8 MiB    122.6 MiB     1407868                   joined_table.insert(s, r)
    55                                             
    56    203.8 MiB      0.0 MiB           1           return joined_table

从profiling结果可直接定位核心问题:仅10万行级别的测试数据,第一次哈希连接就带来122MiB的内存增长,且连接后元组长度从2列膨胀到4列,后续两次连接元组长度会继续增长到6列、8列,叠加Python原生tuple、int对象的大量额外内存开销,1000万行规模下内存占用会轻松突破几十GiB,触发OOM是必然结果。

优化方案

按落地优先级从高到低排序:

  • 替换Python原生存储结构为紧凑列存:原生Python每个int对象占28字节、每个tuple有40字节以上的额外开销,而实际哈希后的整数用4字节就能存储。将PropertyTable的list of tuples存储替换为numpy.ndarray(指定dtype为int32/int64)或pyarrow.Table这类紧凑列存结构,内存占用可直接降至原来的1/10~1/5,优化效果最明显。
  • 连接时裁剪冗余列:当前每次连接会无差别拼接两个表的所有列,比如第一次连接follows和friendOf后,follows.object和friendOf.subject值完全相同,属于重复存储;后续连接仅需保留下一轮要用的连接键和最终查询需要输出的列,无用列直接丢弃,不要带入下一轮连接。
  • 优化哈希表构建逻辑:哈希连接时永远选择行数更少的表构建哈希表,不要固定使用第一个入参建表,尽可能缩小哈希表的内存占用。如果单表规模大到无法放入内存,改用Grace Hash Join:将两个表按连接键的哈希值拆分为多个可完全加载进内存的小分块,逐块完成连接,不需要全量加载数据。
  • 及时回收无用内存:每完成一轮连接,手动删除上一轮中间表的对象引用,触发GC回收内存,避免多轮连接的中间结果同时驻留内存。
  • 调整连接顺序:提前统计4张属性表的行数,按照基数从小到大的顺序安排连接,优先连接数据量最小的表,尽可能控制中间结果的膨胀规模。

额外提示:如果最终查询结果本身规模大到无法放入内存,不要把结果存在内存的PropertyTable对象中,边连接边将结果写入磁盘文件即可。


内容的提问来源于stack exchange,提问作者kklaw

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 05:48:09