【问题标题】:disruptor producer too fast for consumer破坏者生产者对消费者来说太快了
【发布时间】:2015-09-02 19:05:10
【问题描述】:

我正在为某些业务逻辑使用中断器,该中断器发布到另一个处理 IO 的中断器。发布到 IO 中断器的事件可能到达太快,无法构建和验证 IO。嗯,这就是重点......

IO 中断器的设置如下:

disruptor = new Disruptor<>(factory, RING_SIZE, executor, ProducerType.SINGLE, new BlockingWaitStrategy());
disruptor.handleEventsWith(new Logic(disruptor, io));

然后逻辑事件处理程序是这样设置的:

public void onEvent(FixEvent event)
{
     quickfix.Message ioMessage = event.message;
     quickfix.SessionID receiver = event.session;

     Log.debug("message: " + event.message.toString());

     SessionID id = new SessionID(receiver.getBeginString(), "MYFX", receiver.getTargetCompID());
     Session session = Session.lookupSession(id);

     Log.debug("message: " + ioMessage.toString());

     session.send (ioMessage);
}

当您发送 (ioMessage) 时,发生了一个新事件,它以某种方式覆盖了 ioMessage,因此重复的消息被发送出去。

你有什么建议?

【问题讨论】:

  • 我认为您需要添加更多关于 send() 和“做一些工作”部分的详细信息。根据您上面显示的内容,您应该没有问题。 Disruptor 确保您只有一个线程处理 ioMessage,但如果您将该 ioMessage 发送到其他地方,您可能会看到问题。
  • @jasonk 好的,再填写一些。也许问题是日志,在第一个实例中我正在查看 event.message,在第二个实例中查看 ioMessage 变量。这就是我有时会在负载下看到不同的 FIX 事件消息的地方,但问题是有重复的 FIX 消息一个接一个地发送到接收器。当 2 个事件相继发布时,其中的第 2 个事件会以某种方式覆盖第一个事件并被发送两次……很奇怪。

标签: producer-consumer disruptor-pattern


【解决方案1】:

答案看起来像是要完成事件处理程序变量并在事件上使用同步的最终锁定,如下所示:

private final Object lock = new Object();

public void onEvent(final FixEvent event, final long sequence, final boolean endOfBatch)
{   
    synchronized(lock)
    {
        ...

除非锁定是最终的,否则它不起作用。它不适用于静态最终锁。它在除最终锁之外的其他任何东西上同步都不起作用,即在final event 上同步不起作用。

然后它在乐队营地工作了 1 次,然后停止工作......

【讨论】:

  • 中断器已经确保事件处理程序在自己的线程上运行,问题必须存在于堆栈的更下方。
  • @SamTurtelBarker 以及 synchronized(lock) 块内的内容是找到会话 ID 并发送消息,但不知何故消息被覆盖了。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-10-10
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多