【问题标题】:ZMQ Pub Sub - Should it be dropping messages?ZMQ Pub Sub - 它应该丢弃消息吗?
【发布时间】:2014-09-17 18:30:38
【问题描述】:

正在尝试测试 pub/sub 的简单实现。我发现如果我离开订阅者并发送消息,订阅者不会全部收到它们。有时全部收到,有时部分收到,有时整套未收到。

运行订阅服务器(让它一直运行),然后多次运行发布服务器。

import org.zeromq.ZMQ;
import org.zeromq.ZMQ.Context;
import org.zeromq.ZMQ.Socket;


public static void main (String[] args) {

    // Prepare our context and subscriber
    Context context = ZMQ.context(1);
    Socket subscriber = context.socket(ZMQ.SUB);

    subscriber.connect("tcp://localhost:5563");
    subscriber.subscribe("B".getBytes());


    System.out.println("Starting Subscriber..");
    int i = 0;
    while (true) {
        String address = subscriber.recvStr();
        String contents = subscriber.recvStr();
        System.out.println(address+":"+new String(contents) + ": "+(i));
        i++;
    }

}

}

出版商:

import org.zeromq.ZMQ;
import org.zeromq.ZMQ.Context;
import org.zeromq.ZMQ.Socket;


public class TestPublisher {

public static void main (String[] args) throws Exception {
    Context context = ZMQ.context(1);
    Socket publisher = context.socket(ZMQ.PUB);
    publisher.bind("tcp://*:5563");
    System.out.println("Starting Publisher..");
    publisher.setIdentity("B".getBytes());
    publisher.setHWM(1000);
    for (int i = 0; i < 10; i++) {
        Thread.sleep(10l);
        publisher.sendMore("B");
        boolean isSent = publisher.send("We would like to see this:"+i);
        System.out.println("Message was sent "+i+" , "+isSent);
    }

    Thread.sleep(1000);
    publisher.close ();
    context.term ();
}

}

【问题讨论】:

    标签: java publish zeromq subscribe


    【解决方案1】:

    经过一番调试,发现问题是在发布套接字绑定时花费了一些时间,并且尝试发布只是丢弃了消息。在初始绑定上添加一个简单的 100 毫秒睡眠修复了这个问题。在 prod 环境中,发布者在启动时就已经绑定了。

    猜猜这是一个单行解决方案。现在,具有平均数据量的发布/订阅的所有消息都可以正常工作,而不会丢失任何数据。请参阅下面我的发布者的代码 sn-p 更新。

    import org.zeromq.ZMQ;
    import org.zeromq.ZMQ.Context;
    import org.zeromq.ZMQ.Socket;
    
    
    public class TestPublisher {
    
        public static void main (String[] args) throws Exception {
            Context context = ZMQ.context(1);
            Socket publisher = context.socket(ZMQ.PUB);
    
            publisher.bind("tcp://*:5563");
            System.out.println("Starting Publisher..");
            publisher.setIdentity("B".getBytes());
            // for testing setting sleep at 100ms to ensure started.
            Thread.sleep(100l);
            for (int i = 1; i <= 10; i++) {
                publisher.sendMore("B");
                boolean isSent = publisher.send("X("+System.currentTimeMillis()+"):"+i);
                System.out.println("Message was sent "+i+" , "+isSent);
            }
    
            publisher.close ();
            context.term ();
        }
    }
    

    【讨论】:

      【解决方案2】:

      zmq guide 的第 5 章“高级 Pub-Sub 模式”介绍了可靠的 pub/sub。但是,如果您的服务器在生产环境中始终处于启动状态并运行,那么您不需要做任何其他事情,只需让您的测试正常运行。

      如果您真的想解决订阅者丢失消息的问题,第 5 章中的“获取带外快照”示例将涵盖它。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2017-07-12
        • 2011-11-20
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多