KPBOT / CS144 Lab1 笔记 — StreamReassembler:乱序字节流重组

Created Mon, 20 Jul 2026 00:00:00 +0000 Modified Mon, 20 Jul 2026 00:00:00 +0000

经过 Day0 对 Lab0 的完成,今天的任务也很明了:完成 Lab1。

CS144 本质上只实现一个简易 TCP 框架,所以本课程实际涉及到的计算机网络层级一般只用 TCP/IP 协议中的传输层与网络层以及应用层。本次实验也是同理,我们要完成一个 StreamReassembler——将字符串拼接字节流。

它接收:

  • 一段数据 data
  • 该数据在完整字节流中的起始位置 index
  • 是否为结束标志 eof

然后将这些片段按照正确顺序重新组合,输出给后面的 ByteStream。

说实话,今天的任务比较难,涉及到了大量数据结构方面的算法,所以今天也就不过多的回顾计算机网络相关的知识了,主要记录一下做这个 StreamReassembler 的思路。

变量定义

要实现 StreamReassembler 类的功能,首先要了解它实现了什么、用到了什么,然后基于所给的资源把这些资源转换为我们设想要做的东西。

先看最初提供的 stream_reassembler 头文件吧。它预先提供了 2 个私有变量:_capacity_output,以及要我们实现 3 个函数:

  • void push_substring(const string &data, const size_t index, const bool eof)
  • size_t unassembled_bytes()
  • bool empty()

很明显,后面的 unassembled_bytes()empty() 功能都很单一——分别返回接收到但还未发送到 output 的 data 数量(即接收缓存的 length),以及接收缓存是否清空。这两个功能其实都可以通过一个变量来表示其状态:unassembled_bytes() 返回这个变量的值,empty() 返回这个变量是否为空(== 0)。

那么剩下的 push_substring 就变成了我们主要要实现的目标。

可以看到 push_substring 分别接收 3 个参数:data 表示接收到的数据,index 表示这个 data 在总的 data 中位于哪一个起始位置,eof 表示接收到的数据是否结束。

由于发来的 data 十分零碎,要想把这些 data 整合为一个完整的 data,必须进行如下考虑:

  1. 终止判断 —— 因为每一个 data 都是完整 data 的一部分,完整 data 的截止符根据 eof 进行决定。需要一个状态表示当前块的数据结束,即 _eof_flag,避免下一次重复传来的数据被继续写入。

  2. 有序存储 —— 由于发来的 data 并不一定连续,呈现零散分布,所以需要一个存储结构将各个 data 块进行存储,然后合并。为了方便合并,这个结构最好是有序的 index。基于此,可以考虑 set 对这个结构进行存储。

set 底层通常使用红黑树实现,最大的价值在于自动维护元素有序性,使我们能够快速定位当前 block 前后的相邻区间。这里的去重性质并不是依据 data 内容,而是依据 operator< 定义的比较关系——因此通过 begin 作为排序关键字,可以保证同一起始位置的 block 不会重复存在。

但是值得注意的是,就算不考虑 index,我们的 data 是完全顺序写入的,也不能直接将 data 写入 set。因为假设发送 data1=abcdedata2=defgh,直接:

for (int i = 0; i < data.length(); i++) {
    set.add(data[i]);
}

的结果就是 set 内的东西输出出来可能就是 abcdefgh 而不是 abcdedefgh 了。

所以必须得用一个可以存储写入 data 块的起始 index、data 本身、以及结束信息的一个结构体,并且通过实现 bool operator<(const block_node &t) const { return begin < t.begin; } 来让 set 自然地进行从小到大排序。

最终,设计出结构体:

struct block_node {
    size_t begin = 0;
    size_t length = 0;
    std::string data = "";
    bool operator<(const block_node &t) const { return begin < t.begin; }
};

begin 存入这一块数据的 index,length 存储这块数据的长度,从而可以间接判断这个 data 的 end()。

考虑到各个 block 可能存在 A block 存储 [1,10]、B block 存储 [5,15] 的问题,必须让 block 能够进行融合,然后将最终融合后的 block 的数据取到 output。所以创建一个函数 merge_block 来对两个 block 之间进行融合。

同时还要考虑到数据的完整性。head_index 类似 TCP 接收端中的 RCV.NXT,表示当前已经连续接收并交付给上层的最后位置后的下一个字节编号。维护一个 head_index,通过不断写入数据、不断更新头部,就可以确保 head_index 前的数据是完整的。

最终,新增了如下变量与函数:

struct block_node {
    size_t begin = 0;
    size_t length = 0;
    std::string data = "";
    bool operator<(const block_node &t) const { return begin < t.begin; }
};

std::set<block_node> _blocks = {};
size_t _unassembled_byte = 0;
size_t _head_index = 0;
bool _eof_flag = false;

long merge_block(block_node &elm1, const block_node &elm2);

个人感悟:整个过程让我想到了 DNA 的后随链复制——这些零碎的 data 就像一个个冈崎片段,merge_block 是 DNA 连接酶,将序列正确链接;而 push_substring 是主链复制过程中 DNA 聚合酶的位置,不断接收新的碱基,通过配对形成短链 DNA 链条,然后交给 DNA 连接酶进行连接。但也有略微不同:DNA 的复制是一个碱基一个碱基地进行配对,不存在 [1,10] + [5,15] 的重叠现象。

merge_block

merge_block 用于将两个 block 进行融合,并且是有序的融合。基于这么一个目的,需要考虑以下情形与操作:

先对两个 block 进行排序,确定谁的 begin index 在前:

if (elm1.begin > elm2.begin) {
    x = elm2;
    y = elm1;
} else {
    x = elm1;
    y = elm2;
}

通过这种方式简单地对这两个 block 分清前序列与后序列。

然后确认这两个 block 是否连续。 其实很简单,也就是 x 的 index + x.length(即 x 的尾序列号)是否大于等于 y 的 index。只要大于,这两个 block 就可以融合,否则它们孤立:

if (x.begin + x.length < y.begin)
    return -1;

那么假设 x 的 index + x.length >= y.index + y.length 了呢? 这种情况其实就意味着 y 是 x 的子集,所以这个 block 直接就是 x。既然 y 是 x 的子集,那么 y 有多长融入了 x 呢?很显然,直接就是 y 的 length:

if (x.begin + x.length >= y.begin + y.length) {
    elm1 = x;
    return y.length;
}

现在只剩下最后一种情况: 也就是先前提到的 [1,10] + [5,15] 类似的情况。这种情况的操作也很简单——首先这个 block 的 begin index 肯定是 x.begin,其 end index 肯定就是 y.begin + y.length 了。而数据呢?数据当然就是 x.begin 到 x.end 的部分再加上 y.end - x.end 的部分。不过定义里没有定义 block_node 的 end 属性,所以将其化作等价形式,得到最终的合并代码:

elm1.begin = x.begin;
elm1.data = x.data + y.data.substr(x.begin + x.length - y.begin);
elm1.length = elm1.data.length();
return x.begin + x.length - y.begin;

通过以上操作,得到 merge_block 最终完整代码:

long StreamReassembler::merge_block(block_node &elm1, const block_node &elm2) {
    block_node x, y;
    if (elm1.begin > elm2.begin) {
        x = elm2;
        y = elm1;
    } else {
        x = elm1;
        y = elm2;
    }
    if (x.begin + x.length < y.begin) {
        return -1;
    } else if (x.begin + x.length >= y.begin + y.length) {
        elm1 = x;
        return y.length;
    } else {
        elm1.begin = x.begin;
        elm1.data = x.data + y.data.substr(x.begin + x.length - y.begin);
        elm1.length = elm1.data.length();
        return x.begin + x.length - y.begin;
    }
}

push_substring

工具做好了,接下来实现 push_substring

合法性检查

检查 index 是否超过了 _head_index + _capacity

if (index >= _head_index + _capacity) {
    return;
}

定义 elm 并分类处理

定义 block_node elm,然后检查 data 是否已经在之前被接收(即 index + data.length() <= _head_index)。如果已经被接收了,就需要考虑这段数据是否已经到达结尾。若到达结尾则对 _output.end_input()

if (index + data.length() <= _head_index) {
    goto JUDGE_EOF;
}

JUDGE_EOF:
    if (eof) {
        _eof_flag = true;
    }
    if (_eof_flag && empty()) {
        _output.end_input();
    }

处理部分重叠

检查当前片段 data 是否与之前得到的连续序列重合。除去前面 index + data.length() <= _head_index 的情况后,剩下的是 index <= _head_indexindex + data.length() > _head_index 的情况。计算超出前序连续序列部分的偏移量 offset,然后根据 offset 截取对应窗口的 data 到 block 中,方便后续 merge:

else if (index < _head_index) {
    size_t offset = _head_index - index;
    elm.data.assign(data.begin() + offset, data.end());
    elm.begin = _head_index;
    elm.length = elm.data.length();
}

处理全新数据

最后剩余一种情况——data 完全不在连续的前序序列中,直接将完整 data 塞入 block:

else {
    elm.length = data.length();
    elm.begin = index;
    elm.data = data;
}

经过以上操作,_unassembled_byte 自然也就增加了 elm.length

合并状态机

然后开始对各个接收到 block 进行 merge。由于 data 是一个一个传进来的,所以得做成一个状态机,对输入的 data 尽可能进行合并,最终让其融入到整个序列中。

当一个新的 data 进入函数中时,从 set 中寻找一个不大于 elm 的 begin index 的最大 block(即迭代器 iter),然后尝试将其与当前 elm 进行融合。先进行后序融合:当 iter == blocks.end() 时代表当前状态下后序已经没有 block 可融合,所以停止融合;当 merged_byte < 0 时代表没有可以连通当前 block 的 block,也可以停止。

每次融合都在不断消除 blocks 中被融合后的序列。这相当于 iter 从 set 中弹出在当前 block 后的优先队列队首,让 block 与之融合,直到队首为空或无法再融合为止。

然后是让 block 与前面的 block 融合——经过后序融合后 iter 指向的是当前 block,所以需要 iter--,然后开始模拟栈操作:让栈弹出栈顶(依旧是寻找一个不大于当前 elm index 的最大 block),依次融合,直到 blocks 中在当前 block 前的 block 为空,或无法形成连续区段。最终结束这一融合流程,将最终融合结果插入 blocks,等待下一次 data 传入重复此操作。

do {
    long merged_bytes = 0;
    auto iter = _blocks.lower_bound(elm);
    while (iter != _blocks.end()
           && (merged_bytes = merge_block(elm, *iter)) >= 0) {
        _unassembled_byte -= merged_bytes;
        _blocks.erase(iter);
        iter = _blocks.lower_bound(elm);
    }
    if (iter == _blocks.begin()) {
        break;
    }
    iter--;
    while ((merged_bytes = merge_block(elm, *iter)) >= 0) {
        _unassembled_byte -= merged_bytes;
        _blocks.erase(iter);
        iter = _blocks.lower_bound(elm);
        if (iter == _blocks.begin()) {
            break;
        }
        iter--;
    }
} while (false);
_blocks.insert(elm);

交付连续数据到 output

最后将本次操作中得到的、与 _head_index 直接连续融合部分的 data 写入 _output,并更新 _head_index,完成本次数据传入:

if (!_blocks.empty() && _blocks.begin()->begin == _head_index) {
    const block_node head_block = *_blocks.begin();
    size_t write_bytes = _output.write(head_block.data);
    _head_index += write_bytes;
    _unassembled_byte -= write_bytes;
    _blocks.erase(_blocks.begin());
}

完整代码

StreamReassembler::StreamReassembler(const size_t capacity)
    : _output(capacity), _capacity(capacity) {}

long StreamReassembler::merge_block(block_node &elm1, const block_node &elm2) {
    block_node x, y;
    if (elm1.begin > elm2.begin) {
        x = elm2; y = elm1;
    } else {
        x = elm1; y = elm2;
    }
    if (x.begin + x.length < y.begin) {
        return -1;
    } else if (x.begin + x.length >= y.begin + y.length) {
        elm1 = x;
        return y.length;
    } else {
        elm1.begin = x.begin;
        elm1.data = x.data + y.data.substr(x.begin + x.length - y.begin);
        elm1.length = elm1.data.length();
        return x.begin + x.length - y.begin;
    }
}

void StreamReassembler::push_substring(const string &data,
                                       const size_t index,
                                       const bool eof) {
    if (index >= _head_index + _capacity) return;

    block_node elm;
    if (index + data.length() <= _head_index) {
        goto JUDGE_EOF;
    } else if (index < _head_index) {
        size_t offset = _head_index - index;
        elm.data.assign(data.begin() + offset, data.end());
        elm.begin = _head_index;
        elm.length = elm.data.length();
    } else {
        elm.length = data.length();
        elm.begin = index;
        elm.data = data;
    }

    _unassembled_byte += elm.length;
    do {
        long merged_bytes = 0;
        auto iter = _blocks.lower_bound(elm);
        while (iter != _blocks.end()
               && (merged_bytes = merge_block(elm, *iter)) >= 0) {
            _unassembled_byte -= merged_bytes;
            _blocks.erase(iter);
            iter = _blocks.lower_bound(elm);
        }
        if (iter == _blocks.begin()) break;
        iter--;
        while ((merged_bytes = merge_block(elm, *iter)) >= 0) {
            _unassembled_byte -= merged_bytes;
            _blocks.erase(iter);
            iter = _blocks.lower_bound(elm);
            if (iter == _blocks.begin()) break;
            iter--;
        }
    } while (false);
    _blocks.insert(elm);

    if (!_blocks.empty() && _blocks.begin()->begin == _head_index) {
        const block_node head_block = *_blocks.begin();
        size_t write_bytes = _output.write(head_block.data);
        _head_index += write_bytes;
        _unassembled_byte -= write_bytes;
        _blocks.erase(_blocks.begin());
    }

JUDGE_EOF:
    if (eof) _eof_flag = true;
    if (_eof_flag && empty()) _output.end_input();
}

size_t StreamReassembler::unassembled_bytes() const { return _unassembled_byte; }

bool StreamReassembler::empty() const { return _unassembled_byte == 0; }