Tweaking the producer-consumer model

·

The classic producer-consumer pattern makes a few assumptions in order to work. In this post I’ll discuss some of these assumptions, what happens when they break, and a cool solution to deal with it. The two assumptions I’d like to discuss are:
  1. All producers are equal - this means that all producers are limited by one, shared limit - the mediator queue between them and the consumers. This is a great assumption for most implementations: usually you are limited by the amount of work your consumers can handle and usually those are limited by the amount of processors you would physically have, therefore it doesn’t matter if the piece of work came from producer X or producer Y - either way you’d still have the same amount of processors.
  2. Once sent to the queue, a producer doesn’t care about its products - this means that there is no link between the produced work and the producer once it has been sent to the mediation queue. And why should it? Its task was to create the work, afterwards its the consumer’s job to deal with the work, and then it would be the garbage collector's task to reclaimed the memory space that work used. If the producer kept a link to the work, it couldn’t have been reclaimed without proper notifications between the producer and the consumer - and that would be completely against the decoupling nature of the producer-consumer pattern.
So far you might be thinking “these are good, based assumptions; why should they ever break?”. Well, keep on reading then.

When good, based assumptions break

Why would these assumptions break? Usually, they wouldn’t. However, there are certain end-cases where they might. I’ll follow this break with an example to illustrate it better.
  1. Different producers should have different limits - this can occur when you know that certain producers create resource-consuming tasks, such as requiring a lot of a consumer’s physical memory. If that’s the case, it might be wise to limit those producers separately. The naive way to do so would be to create another mediator queue; however, such an implementation could be terrible to manage since it would require the producers to poll both queues in a loop, possibly wasting CPU cycles where before they would just be blocking until new work would arrive.
  2. Producers should care about their products’ end-of-life - following the previous example, this is a natural solution to it: if the producer knew when its products are out of the consumers’ memory, they could then send the next work item - keeping the consumer’s memory usage as low as required.
We could phrase our required solution as “what we want is to limit the number of products coming from a certain producer so that only a certain number of them are being dealt with by the consumers at any given time”.

The solution could be simple, if someone just listened

The easiest way to achieve this would be to use a semaphore within the producer, acquiring a lock whenever a new product is created and releasing a lock whenever that product is completely consumed. This solution can be somewhat decoupled too: piggie-back a listener on the product, and make the consumer call the listener’s invocation when the consumption is finished: Listener approach That means that the producer and consumer need to use a wrapper for the work unit, since it needs to contain the listener as well as the raw data. For that, we could create a Work class:

final class Work<T> {
  public final T data;
  public final ConsumedListener listener;

  public <U> Work<U> spawn(U data) {
    // create a new work unit with the same listener, this is for map/reduce cases
  }

  // ctor
}
There are a few problems with this kind of solution though, mostly because it couples the consumer with the producer’s logic a bit. The consumer is forced to “release” the product using the listener - something it wouldn’t do for other products. This could get difficult when you realise you need to think of complex, multi-layer consumers, exception handling and consumers that use map/reduce and split the work, which can in turn cause the consumer to accidentally call the listener more than once. What comes to mind is “if we only knew when the work is no longer used by the consumer automatically”. Well, we actually do know: this is where some garbage collection tricks can come in. For example, we could use the finalize method to call that listener automatically! That would handle exception cases, the splitting cases, and the complex consumer cases, while relieving any work from the consumer itself - it wouldn’t know its dealing with a different type of product than the rest.

Back to some old GC tricks

If you have been following the garbage collection series, you probably saw the post about garbage collection tips. One of the tips was about why you should avoid using the finalize method and how to do so efficiently. If you haven’t, I recommend you do so now. It could make the rest of the post much easier to understand. So, knowing we want to use the GC to tell us when the product is no longer referenced, and at the same time not wanting to use the finalize method, we’re left with one solution: weak(er) references. In fact, I’d use PhantomReference in this code, even though I could probably use any of the others; the PhantomReference is just guaranteed never to keep the object from being reclaimed, no matter what. The solution is as follows: when the producer creates the product, it acquires a lock from the semaphore. If getting a lock was successful, it creates a reference object, and wraps both the product and the reference object together in the previously mentioned, slightly altered Work wrapper object. The producer also creates a phantom reference to the reference object, and registers it with a reference queue. If getting a lock was unsuccessful, it waits a certain amount of time, and checks the reference queue: every phantom reference there marks a reference object claimed by the GC, so it releases a lock, and trying to acquire the lock again (in a loop). This might be better illustrated with a drawing: Phantom approach

Conclusions?

While the producer-consumer pattern offers a great solution for work distribution, either in-bounds or out-of-bounds, with almost no development efforts, sometimes it does require some tweaking to get it right. It might be that the solution I offer above is not the best for my own case; maybe a different pattern can be used? I’d be more than grateful to hear your experience with tweaking the producer-consumer model, or your ideas and thoughts about the implementation above. And before I forget: implementation files for the Work and PhantomProducer classes, and the test class!

Comments (3)

Aviad
@Duski: Yes, you're correct. The assumption the model I've suggested is that there's a constant flow of information through the system; therefore, a GC call is inevitably called soon enough. In my tests, I didn't have to wait long for a GC to occur (check the unit-tests where I simulated this behavior by creating a certain amount of objects and releasing them just to stir up the GC to take action). In an environment where not many objects are created and released, a less preferred approach can be taken which is to call System.gc() within the phantom producer's inner loop. Also take note that this approach is not good for high throughput producers as these are inherently slower due to the timeout block in the inner loop. For high throughput producers, a different approach should probably be used.
Duski
Doesn't this solution rely on the fact that GS will kick in soon after the unit of work get's out of scope in consumer? If that doesn't happen, your processing will be sitting idle or you will be not using resources that you have effectively.
Duski
Sorry for the typo, with GS I meant GC.