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

如何在Dask中实现100GB有序区块链数据集的合并?

问题

我有一个约100GB的区块链数据集,需要基于transactionHash字段合并两个数据表。由于两张表都已经按blockNumber字段排序,理论上可以通过O(|A|+|B|)=O(n)的低复杂度完成合并(而非常规的O(n²)复杂度)。

在Pandas中可以用pd.merge_ordered实现这个操作,但Dask里没有对应的功能。

附示例代码:

blockNumber = [1,1,1, 2,2,2]
transactionHash = ['0x0', '0x1', '0x6', '0x2', '0x8', '0xf']

df1 = pd.DataFrame({
    'blockNumber': blockNumber,
    'transactionHash': transactionHash,
    'field1': [42,37,21,78,32,45],
})
df2 = pd.DataFrame({
    'blockNumber': blockNumber,
    'transactionHash': transactionHash,
    'field2': [True, False, True, True, False, False],
})
res1 = pd.merge(df1, df2, on=['transactionHash'])
res2 = pd.merge(df1, df2, on=['blockNumber', 'transactionHash'])
assert res1.equals(res2)

我想了解:
a) 自行实现该算法需考虑哪些要点,如何对齐分区;
b) 可使用哪些Dask替代函数完成合并。


解决方案

a) 自行实现算法的要点与分区对齐

核心算法要点

  • 双指针遍历逻辑:两张表均按blockNumber排序,且同区块内transactionHash唯一有序(区块链交易特性),可采用双指针分别遍历两表行:
    • 若两指针指向的blockNumber相等,对比transactionHash,匹配则合并行;若A的哈希更小则移动A指针,反之移动B指针。
    • 若A的blockNumber小于B的,移动A指针;反之移动B指针。
  • 重复键兼容:虽然交易哈希通常唯一,但仍需处理单表内哈希重复的场景,支持一对多/多对多匹配逻辑。
  • 边界处理:遍历到表末尾时终止循环,处理其中一张表先遍历完成的剩余数据。

分区对齐方案

Dask并行处理的核心是分区对齐,需保证分区对应连续的blockNumber区间:

  1. 统一分区边界:先获取df1各分区的blockNumber最小/最大值,让df2按照相同边界重新分区(用dd.repartition指定自定义divisions)。
  2. 分区级并行合并:对应分区内的数据可独立执行双指针合并,因为同分区的blockNumber连续且不跨区,不会出现跨分区匹配项。
  3. 边界区块补全:若存在blockNumber落在分区边界的情况,提取每个分区的边界区块(如最后一个区块的所有数据),与相邻分区的边界数据合并处理后,再合并回最终结果,避免遗漏匹配。

b) Dask替代合并函数

1. map_partitions结合Pandas有序合并

利用map_partitions对对齐后的分区,调用Pandas的merge_ordered或自定义双指针函数处理:

import dask.dataframe as dd
import pandas as pd

# 先对齐df1和df2的分区(确保divisions完全一致)
df2 = df2.repartition(divisions=df1.divisions)

def merge_single_partition(df1_part, df2_part):
    return pd.merge_ordered(
        df1_part, df2_part,
        on=['blockNumber', 'transactionHash'],
        fill_method=None
    )

# 定义结果元数据
meta = pd.concat([df1.head(1), df2.head(1)], axis=1).drop_duplicates()
merged_df = dd.map_partitions(merge_single_partition, df1, df2, meta=meta)

2. dd.merge配合预分区优化

提前按blockNumber分区,减少全局shuffle,让Dask仅在同区块分区内做哈希合并:

# 按blockNumber设置索引并排序,确保分区对齐
df1 = df1.set_index('blockNumber', sorted=True)
df2 = df2.set_index('blockNumber', sorted=True)

# 合并时利用索引分区减少数据交换
merged_df = dd.merge(
    df1, df2,
    left_index=True, right_index=True,
    on='transactionHash'
)

3. merge_asof(适合严格有序场景)

若transactionHash在同blockNumber内严格递增,可使用专门的有序合并函数merge_asof,复杂度为O(n):

merged_df = dd.merge_asof(
    df1.sort_values(['blockNumber', 'transactionHash']),
    df2.sort_values(['blockNumber', 'transactionHash']),
    on='blockNumber',
    by='transactionHash'
)

内容的提问来源于stack exchange,提问作者David Davíd

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 06:43:11