【问题标题】:ZeroMQ, can we use inproc: transport along with pub/sub messaging patternZeroMQ,我们可以使用inproc:传输以及发布/订阅消息模式吗
【发布时间】:2016-02-17 10:18:46
【问题描述】:

场景:

我们正在评估 ZeroMQ(特别是 jeroMq)的事件驱动机制。

应用程序是分布式的,其中多个服务(发布者和订阅者都是服务)可以存在于同一个 jvm 或不同节点中,这取决于部署架构。

观察

为了玩耍,我创建了一个 pub/sub 模式,使用 inproc: 作为传输方式,使用 jero mq(版本:0.3.5)

  1. 线程发布能够发布(貌似发布了,至少没有错误)
  2. 另一个线程中的订阅者没有收到任何东西。

问题

使用 inproc:pub/sub 可行吗?

尝试了谷歌搜索,但找不到任何具体的信息或见解?

pub/subinproc: 的代码示例

使用 jero mq(版本:0.3.5)的 inproc pub sub 的工作代码示例对以后访问这篇文章的人很有用。一位发布者发布主题 AB,两位订阅者分别接收 AB

/**
 * @param args
 */
public static void main(String[] args) {

    // The single ZMQ instance
    final Context context = ZMQ.context(1);

    ExecutorService executorService = Executors.newFixedThreadPool(3);
    //Publisher
    executorService.execute(new Runnable() {

        @Override
        public void run() {
            startPublishing(context);
        }
    });
    //Subscriber for topic "A"
    executorService.execute(new Runnable() {

        @Override
        public void run() {
            startFirstSubscriber(context);
        }
    });
    // Subscriber for topic "B"
    executorService.execute(new Runnable() {

        @Override
        public void run() {
            startSecondSubscriber(context);
        }
    });

}

/**
 * Prepare the publisher and publish
 * 
 * @param context
 */
private static void startPublishing(Context context) {

    Socket publisher = context.socket(ZMQ.PUB);
    publisher.bind("inproc://test");
    while (!Thread.currentThread().isInterrupted()) {
        // Write two messages, each with an envelope and content
        try {
            publisher.sendMore("A");
            publisher.send("We don't want to see this");
            LockSupport.parkNanos(1000);
            publisher.sendMore("B");
            publisher.send("We would like to see this");
        } catch (Throwable e) {

            e.printStackTrace();
        }
    }
    publisher.close();
    context.term();
}

/**
 * Prepare and receive through the subscriber
 * 
 * @param context
 */
private static void startFirstSubscriber(Context context) {

    Socket subscriber = context.socket(ZMQ.SUB);

    subscriber.connect("inproc://test");

    subscriber.subscribe("B".getBytes());
    while (!Thread.currentThread().isInterrupted()) {
        // Read envelope with address
        String address = subscriber.recvStr();
        // Read message contents
        String contents = subscriber.recvStr();
        System.out.println("Subscriber1 " + address + " : " + contents);
    }
    subscriber.close();
    context.term();

}

/**
 * Prepare and receive though the subscriber
 * 
 * @param context
 */
private static void startSecondSubscriber(Context context) {
    // Prepare our context and subscriber

    Socket subscriber = context.socket(ZMQ.SUB);

    subscriber.connect("inproc://test");
    subscriber.subscribe("A".getBytes());
    while (!Thread.currentThread().isInterrupted()) {
        // Read envelope with address
        String address = subscriber.recvStr();
        // Read message contents
        String contents = subscriber.recvStr();
        System.out.println("Subscriber2 " + address + " : " + contents);
    }
    subscriber.close();
    context.term();

}

【问题讨论】:

  • 添加场景示例代码以备参考

标签: java zeromq publish-subscribe event-driven-design jeromq


【解决方案1】:

The ZMQ inproc transport 旨在用于不同线程之间的单个进程中。当您说“可以存在于同一个 jvm 或不同的节点中”(强调我的)时,我假设您的意思是您将多个进程作为分布式服务而不是单个进程中的多个线程.

如果是这样,那么不,您尝试执行的操作不适用于 inprocPUB-SUB/inproc 可以在多个线程之间的单个进程中正常工作。


编辑以解决 cmets 中的更多问题:

使用inprocipc 之类的传输的原因是,当您在正确的上下文中使用它们时,它比tcp 传输更高效(更快)。可以想象,您可以混合使用多种传输方式,但您始终必须在同一传输方式上绑定和连接才能使其正常工作。

这意味着每个节点最多需要三个PUBSUB 套接字 - 一个tcp 发布者与远程主机上的节点通信,一个ipc 发布者与同一进程上不同进程上的节点通信主机和inproc 发布者在同一进程中与不同线程中的节点通信。

实际上,在大多数情况下,您只需使用 tcp 传输,并且只为所有内容启动一个套接字 - tcp 在任何地方都可以使用。如果每个套接字负责特定的种类信息,则可能启动多个套接字。

如果您总是向其他线程发送一种消息类型而向其他主机发送另一种消息类型是有原因的,那么多个套接字是有意义的,但在您的情况下,从一个节点的角度来看,所有其他节点都是平等的。在那种情况下,我会在任何地方使用tcp 并完成它。

【讨论】:

  • 谢谢@jason,应该是我的示例代码有问题,我会研究一下并放在这里。
  • 另外,在我们的例子中,我们有订阅者/发布者正在处理和其他进程中。关于如何在这种情况下继续前进的任何提示。
  • 如果您要在多个不同的进程之间进行,但都在同一个逻辑服务器实例上,那么ipc 传输是您的最佳选择。如果您要通过网络与其他计算机通信,请使用 tcp。
  • 知道了,inproc 用于线程之间,ipc 用于本地进程之间,tcp 用于远程进程之间。在我们的场景中,可能有发布者/订阅者想要绑定到一个主题,比如 nasdaq 事件。这些发布者/订阅者可以在同一个进程中,也可以在多个进程中(可以在远程和同一个服务器实例中)。我们可以从不同的传输中发布/订阅一个主题吗?
  • 谢谢@Jason,我得到了一张好照片,会在这里更新示例。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-06-24
  • 2016-07-03
  • 2016-11-25
  • 1970-01-01
相关资源
最近更新 更多