当前位置:首页 > PHP教程 > PHP总结归纳

redis源代码分析20–发布/订阅

redis的发布/订阅(publish/subscribe)功能类似于传统的消息路由功能,发布者发布消息,订阅者接收消息,沟通发布者和订阅者之间的桥梁是订阅的channel或者pattern。发布者向指定的publish或者pattern发布消息,订阅者阻塞在订阅的channel或者pattern。可以

redis的发布/订阅(publish/subscribe)功能类似于传统的消息路由功能,发布者发布消息,订阅者接收消息,沟通发布者和订阅者之间的桥梁是订阅的channel或者pattern。发布者向指定的publish或者pattern发布消息,订阅者阻塞在订阅的channel或者pattern。可以看到,发布者不会指定哪个订阅者才能接收消息,订阅者也无法只接收特定发布者的消息。这种订阅者和发布者之间的关系是松耦合的,订阅者不知道是谁发布的消息,发布者也不知道谁会接收消息。

redis的发布/订阅功能主要通过subscribe、unsubscribe、psubscribe、punsubscribe 、publish五个命令来表现。其中subscribe、unsubscribe用于订阅或者取消订阅channel,而psubscribe、punsubscribe用于订阅或者取消订阅pattern,发布消息则通过publish命令。

对于发布/订阅功能的实现,我们先来看看几个与此相关的结构。

struct redisserver {
    ---
   /* pubsub */
   dict *pubsub_channels;/* map channels to list of subscribed clients */
   list *pubsub_patterns;/* a list of pubsub_patterns */
   ---
}
typedef struct redisclient {
   ---
   dict *pubsub_channels; /* channels a client is interested in (subscribe) */
   list *pubsub_patterns; /* patterns a client is interested in (subscribe) */
} redisclient;

在redis的全局server变量(redisserver类型)中,channel和订阅者之间的关系用字典pubsub_channels来保存,特定channel和所有订阅者组成的链表构成pubsub_channels字典中的一项,即字典中的每一项可表示为(channel,订阅者链表);pattern和订阅者之间的关系用链表pubsub_patterns来保存,链表中的每一项可表示成(pattern,redisclient)组成的字典。

在特定订阅者redisclient的结构中,pubsub_channels保存着它所订阅的channel的字典,而订阅的模式则保存在链表pubsub_patterns中。

从上面的解释,我们再来看看订阅/发布命令的最坏时间复杂度(注意字典增删查改一项的复杂度为o(1),而链表的查删复杂度为o(n),从链表尾部增加一项的复杂度为o(1))。

subscribe:

订阅者用subscribe订阅特定channel,这需要在订阅者的redisclient结构中的pubsub_channels增加一项(复杂度为 o(1)),然后在redisserver 的pubsub_channels找到该channel(复杂度为o(1)),并在该channel的订阅者链表的尾部增加一项(复杂度为o(1),注意,如果pubsub_channels中没找到该channel,则插入的复杂度也同为o(1)),因此订阅者用subscribe订阅特定 channel的最坏时间复杂度为o(1)。

unsubscribe:

订阅者取消订阅时,需要先从订阅者的redisclient结构中的pubsub_channels删除一项(复杂度为o(1)),然后在 redisserver 的pubsub_channels找到该channel(复杂度为o(1)),然后在channel的订阅者链表中删除该订阅者(复杂度为o(1)),因此总的复杂度为o(n),n为特定channel的订阅者数。

psubscribe:

订阅者用psubscribe订阅pattern时,需要先在redisclient结构中的pubsub_patterns先查找是否已存在该 pattern(复杂度为o(n)),并在不存在的情况下往redisclient结构中的pubsub_patterns和redisserver结构中的pubsub_patterns链表尾部各增加一项(复杂度都为o(1)),因此,总的复杂度为o(n),其中n为订阅者已订阅的模式。

punsubscribe:

订阅者用punsubscribe取消对pattern的订阅时,需要先在redisclient结构中的pubsub_patterns链表中删除该 pattern(复杂度为o(n)),并在redisserver结构中的pubsub_patterns链表中删除订阅者和pattern组成的映射(复杂度为o(m),因此,总的复杂度为o(n+m),其中n为订阅者已订阅的模式,而m为系统中所有订阅者和所有pattern组成的映射数。

publish:

发布消息时,只会向特定channel发布,但该channel可能会匹配某个pattern。因此,需要先在redisserver结构中的 pubsub_channels找到该channel的订阅者链表(o(1)),然后发送给所有订阅者(复杂度为o(n)),然后查看 redisserver结构中的pubsub_patterns链表中的所有项,看channel是否和该项中的pattern匹配(复杂度为o(m))(注意,这并不包括模式匹配的复杂度),因此,总的复杂度为o(n+m),。其中n为该channel的订阅者数,而m为系统中所有订阅者和所有 pattern组成的映射数。另外,从这也可以看出,一个订阅者是可能多次收到同一个消息的。

解释了发布/订阅的算法后,其代码就好理解了,这里仅给出publish命令的处理函数publishcommand的代码,更多相关命令的代码请参看redis的源代码。

static void publishcommand(redisclient *c) {
    int receivers = pubsubpublishmessage(c->argv[1],c->argv[2]);
    addreplylonglong(c,receivers);
}
/* publish a message */
static int pubsubpublishmessage(robj *channel, robj *message) {
    int receivers = 0;
    struct dictentry *de;
    listnode *ln;
    listiter li;
    /* send to clients listening for that channel */
    de = dictfind(server.pubsub_channels,channel);
    if (de) {
        list *list = dictgetentryval(de);
        listnode *ln;
        listiter li;
        listrewind(list,&li);
        while ((ln = listnext(&li)) != null) {
            redisclient *c = ln->value;
            addreply(c,shared.mbulk3);
            addreply(c,shared.messagebulk);
            addreplybulk(c,channel);
            addreplybulk(c,message);
            receivers++;
        }
    }
    /* send to clients listening to matching channels */
    if (listlength(server.pubsub_patterns)) {
        listrewind(server.pubsub_patterns,&li);
        channel = getdecodedobject(channel);
        while ((ln = listnext(&li)) != null) {
            pubsubpattern *pat = ln->value;
            if (stringmatchlen((char*)pat->pattern->ptr,
                                sdslen(pat->pattern->ptr),
                                (char*)channel->ptr,
                                sdslen(channel->ptr),0)) {
                addreply(pat->client,shared.mbulk4);
                addreply(pat->client,shared.pmessagebulk);
                addreplybulk(pat->client,pat->pattern);
                addreplybulk(pat->client,channel);
                addreplybulk(pat->client,message);
                receivers++;
            }
        }
        decrrefcount(channel);
    }
    return receivers;
}

最后提醒一下,处于发布/订阅模式的client,是无法发布上述五种命令之外的命令(quit除外),这是在processcommand函数中检查的,可以参看前面命令处理章节对该函数的解释。


【说明】本文章由站长整理发布,文章内容不代表本站观点,如文中有侵权行为,请与本站客服联系(QQ:)!