【问题标题】:ZMQ missing events being propagated in jeromq scalaZMQ 丢失事件在 jeromq scala 中传播
【发布时间】:2018-04-22 19:53:11
【问题描述】:

我是 ZeroMQ 的新手,我的 begin() 方法似乎在循环中丢失消息。

我想知道我是否遗漏了一个我没有排队消息的部分?

当我在我的发布者上引发一个事件时,它向我的订阅者发送了两条消息,中间有一个小间隔,我似乎没有收到第二条被中继的消息。 我错过了什么?

class ZMQSubscriber[T <: Transaction, B <: Block](
  socket: InetSocketAddress,
  hashTxListener: Option[HashDigest => Future[Unit]],
  hashBlockListener: Option[HashDigest => Future[Unit]],
  rawTxListener: Option[Transaction => Future[Unit]],
  rawBlockListener: Option[Block => Future[Unit]]) {
  private val logger = BitcoinSLogger.logger

  def begin()(implicit ec: ExecutionContext) = {
    val context = ZMQ.context(1)

    //  First, connect our subscriber socket
    val subscriber = context.socket(ZMQ.SUB)
    val uri = socket.getHostString + ":" + socket.getPort

    //subscribe to the appropriate feed
    hashTxListener.map { _ =>
      subscriber.subscribe(HashTx.topic.getBytes(ZMQ.CHARSET))
      logger.debug("subscribed to the transaction hashes from zmq")
    }

    rawTxListener.map { _ =>
      subscriber.subscribe(RawTx.topic.getBytes(ZMQ.CHARSET))
      logger.debug("subscribed to raw transactions from zmq")
    }

    hashBlockListener.map { _ =>
      subscriber.subscribe(HashBlock.topic.getBytes(ZMQ.CHARSET))
      logger.debug("subscribed to the hashblock stream from zmq")
    }

    rawBlockListener.map { _ =>
      subscriber.subscribe(RawBlock.topic.getBytes(ZMQ.CHARSET))
      logger.debug("subscribed to raw block")
    }

    subscriber.connect(uri)
    subscriber.setRcvHWM(0)
    logger.info("Connection to zmq client successful")

    while (true) {
      val notificationTypeStr = subscriber.recvStr(ZMQ.DONTWAIT)
      val body = subscriber.recv(ZMQ.DONTWAIT)
      Future(processMsg(notificationTypeStr, body))
    }
  }

  private def processMsg(topic: String, body: Seq[Byte])(implicit ec: ExecutionContext): Future[Unit] = Future {

    val notification = ZMQNotification.fromString(topic)
    val res: Option[Future[Unit]] = notification.flatMap {
      case HashTx =>
        hashTxListener.map { f =>
          val hash = Future(DoubleSha256Digest.fromBytes(body))
          hash.flatMap(f(_))
        }
      case RawTx =>
        rawTxListener.map { f =>
          val tx = Future(Transaction.fromBytes(body))
          tx.flatMap(f(_))
        }
      case HashBlock =>
        hashBlockListener.map { f =>
          val hash = Future(DoubleSha256Digest.fromBytes(body))
          hash.flatMap(f(_))
        }
      case RawBlock =>
        rawBlockListener.map { f =>
          val block = Future(Block.fromBytes(body))
          block.flatMap(f(_))
        }
    }
  }
}

【问题讨论】:

    标签: java scala zeromq distributed-system jeromq


    【解决方案1】:

    所以这似乎已经通过在while-loop 中使用ZMsg.recvMsg() 而不是

      val notificationTypeStr = subscriber.recvStr(ZMQ.DONTWAIT)
      val body = subscriber.recv(ZMQ.DONTWAIT)
    

    我不确定为什么会这样,但确实如此。所以这就是我的begin 方法现在的样子

        while (run) {
          val zmsg = ZMsg.recvMsg(subscriber)
          val notificationTypeStr = zmsg.pop().getString(ZMQ.CHARSET)
          val body = zmsg.pop().getData
          Future(processMsg(notificationTypeStr, body))
        }
        Future.successful(Unit)
      }
    

    【讨论】:

      【解决方案2】:

      我错过了什么?

      阻塞 v/s 非阻塞作案手法如何工作:

      诀窍在于对 .recv() 方法的相应调用的(非)阻塞模式

      subscriber.recv( ZMQ.DONTWAIT )-方法的第二次调用因此会立即返回,因此您的第二部分(body)可能并且将在法律上不包含任何内容,即使您的承诺声明确实从发布者发送了一对消息-side (一对.send() 方法调用 - 一个也可能反对,发件人实际上可能只发送一条消息,以多部分方式 - MCVE 代码在这部分不具体)。

      因此,一旦您将代码从非阻塞模式(在 O/P 中)移动到主要阻塞模式(锁定/同步代码的进一步流与到达的外部事件任何格式合理的消息,而不是更早返回),在:

      val zmsg = ZMsg.recvMsg(subscriber) // which BLOCKS-till-a-1st-zmsg-arrived
      

      进一步处理的 .pop() 部分都只是卸载组件(参考实际由发布方发送的ZMsg 多部分结构的备注,如上所示)


      Safety next :
      unlimited alloc -s v/s 强制阻塞/丢弃消息?

      代码在几点让我感到惊讶。除了对.connect()-method 的相当“延迟”调用之外,与之前所有的socket-archetype 详细设置(通常在请求建立连接之后“安排”)相比。虽然这可能按预期工作得很好,但它为.Context()-instance 提供了更紧(更小)的时间窗口来设置和(重新)协商所有相关的连接细节,以便成为 RTO。

      有一条特别的线引起了我的注意:subscriber.setRcvHWM( 0 ) 这是一个依赖于版本原型的技巧。然而,零值会导致应用程序变得易受攻击,我不建议在任何生产级应用程序中这样做。

      【讨论】:

        猜你喜欢
        • 2016-07-11
        • 1970-01-01
        • 1970-01-01
        • 2023-03-09
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2023-02-09
        • 1970-01-01
        相关资源
        最近更新 更多