Go Concurrency Patterns: Pipelines and cancellation
blog.golang.org
blog.golang.org
func merge(cs ...<-chan int) <-chan int {
Inside the function you cannot send to a member of cs - it will give you a compile-error.Outside the function you cannot send to its return value.
An example to play with: http://play.golang.org/p/6AJLpZkLBe
Players send their move to the game on a channel they expose through Move(). They receive regular updates of the game state on a channel they expose through Update().
type Player interface {
Name() string
Move() <-chan Move
Update() chan<- State
}
> https://github.com/aybabtme/bomberman/blob/master/player/pla...The game reads the player moves, if moves are available[1]:
for pState, player := range g.Players {
if pState.Alive {
select {
case m := <-player.Move():
movePlayer(g, board, pState, m)
default:
}
}
}
> https://github.com/aybabtme/bomberman/blob/master/bomberman....The game sends each turn's update to players that are ready to listen for it:
for pState, player := range game.Players {
pState.Board = board.Clone()
pState.Turn = game.Turn()
select {
case player.Update() <- *pState:
default:
}
}
> https://github.com/aybabtme/bomberman/blob/master/bomberman....This leans to surprisingly clean implementations:
* A websocket player [2]
* A keyboard player [3]
* A random AI player [4]
Note that the rest of the game (`bomberman.go`) is a hairy ball that needs refactoring.[1]: this, to prevent an unresponsive player from hanging the game, or an AI player from unfairly computing longer than it's opponent
[2]: https://github.com/aybabtme/bomberweb/blob/master/player/pla...
[3]: https://github.com/aybabtme/bomberman/blob/master/player/inp...
[4]: https://github.com/aybabtme/bomberman/blob/master/player/ai/...
for n := range c {
select {
case out <- n:
case <-done:
}
}
This, is brilliant.Now I can see it allows you do "discard" all the pending reads from the input channel (and avoid writing to a possibly-closed 'out' channel), but wouldn't it be better to break the loop in this case for an immediate exit?
i.e. go for 'interrupt' semantics, rather than 'drain'?
for n := range in {
select {
case out <- n * n:
case <-done:
return
}
}
If you allow pipeline stages to interrupt their receive loop instead of draining the inbound channel, then _all_ send operations need to be governed by done. Otherwise the interrupted receive loop may block the upstream sender.If `break` is used, for loop is terminated once it's "done", hence nothing would be reading from channel `c`. The previous station in the pipeline would then block at sending operation.
ExecutorService pipeline = ExecutorService.newSingleThreadExecutor();
FutureTask<String> step1 =
new FutureTask<String>(new Callable<String>() {
public String call() {
return "step1";
}});
executor.execute(step1); // execute single-step pipeline
step1.cancel(true); // interrupt the taskObviously there is a massive syntactic difference between your Java example and the equivalent Go code, but the other major difference that jumps out at me is that cancellation is baked into the Java library, whereas in Go the programmer determines how each of their pipeline stages should be shut down.
With Java 8 and lambdas, syntax will be much cleaner, i.e. you can write something like: executor.execute(() -> "step1");
Furthermore, is there a way in "Java SingleThreadExecutor" and its related components that allow automatic scheduling tasks on different threads? For example, you are on a quad-core and you want to have 4 threads running together, with millions of tasks mapped to 4 threads and dynamically balanced.
http://docs.oracle.com/javase/7/docs/api/java/util/concurren...
Afaik, what Go does is just an implementation of http://en.wikipedia.org/wiki/Communicating_sequential_proces... which is also available in some other languages (Clojure is the one that comes to mind for me). In Java, I think a closer analogy is a http://docs.oracle.com/javase/7/docs/api/java/util/concurren... which lets the producer and consumer threads automatically work in lock step as well.
In java you also can do something like:
List<Future<Long>> futures = new ArrayList<>();
for(int i = 0; i < n; i ++) futures.add(some parallel execution logic);
long count = 0;
for(Future<Long> f: futures) count += f.get();
These slides go into some more examples: http://talks.golang.org/2013/advconc.slide
You end up with a really nice "actor model", where you have a goro which owns some state responding to incoming messages over channels and sending replies.
Slide 24 gives the general structure, but it's a good read.
The latter has the advantage you get to recompose the system in new ways, yielding greater flexibility.
In Java, a better compositional framework is probably Erik Meijer's et.al's RxJava. Haskell also provides better compositionality due to it's functional structure and lazy evaluation semantics.
For those that would rather read this in presentation form: http://talks.golang.org/2012/concurrency.slide#1
And for those that would rather watch a video presentation: https://www.youtube.com/watch?v=f6kdp27TYZs
http://swtch.com/~rsc/thread/squint.pdf
(it's also linked at the bottom of the article)
check out the "power series, power serious" paper by doug-macllroy. he does that...