I have re-written the core functionality of the data store mentioned in The Reliability Of Go and was surprised as to how few lines of code it boiled down to using the latest version of go.
It's here quickbeam.
Showing posts with label golang. Show all posts
Showing posts with label golang. Show all posts
Saturday, 7 June 2014
Wednesday, 6 November 2013
Dynamic Input Channels
Suppose a function produces some output on a channel, and that the output of this function is to be conditionally consumed dependant on the current system state. Also, lets suppose, the system contains a restriction in that the output channel of the producer function must remain fixed.
The question is then, how to route the output to the correct consuming logic, without reassigning the output channel.
One solution is to dynamically assign the consumer's input channel. Go provides a mechanism for this; the channel itself.
s := make(chan chan int)
c := make(chan int)
for {
select {
case c = <- s:
case v := <- c:
logic(v)
}
}
With this logic in place the consumer's input channel can be dynamically reassigned whilst the producer remains agnostic.
Here is a complete example of the mechanism at work.
Tuesday, 24 September 2013
Bullhorn
Bullhorn is a golang package that provides lightweight type-agnostic publish / subscribe messaging to goroutines.
You can find it here.
The model has been extracted from various time sensitive messaging applications I have worked with over the past few years.
Here is an example limit order / price matching procedure written using bullhorn subscriptions and events.
Running the code, will give you something like...
Got order:{CBG 0 1.05 100}
Got order:{CBG 1 1.05 200}
Matching:{CBG 0 1.05 100} against:CBG: 1655 0.8704 1.4305 2797
Matching:{CBG 1 1.05 200} against:CBG: 1655 0.8704 1.4305 2797
Matching:{CBG 0 1.05 100} against:CBG: 9414 0.8847 1.9453 9074
Matching:{CBG 1 1.05 200} against:CBG: 9414 0.8847 1.9453 9074
Matching:{CBG 0 1.05 100} against:CBG: 8787 0.9490 1.0373 12530 - Execute Order!
Matching:{CBG 1 1.05 200} against:CBG: 8787 0.9490 1.0373 12530
Matching:{CBG 1 1.05 200} against:CBG: 5840 0.9519 1.5722 12054
Matching:{CBG 1 1.05 200} against:CBG: 2485 0.9425 1.7220 11042
Matching:{CBG 1 1.05 200} against:CBG: 5235 1.0583 1.5840 4333 - Execute Order!
Here is an example limit order / price matching procedure written using bullhorn subscriptions and events.
Running the code, will give you something like...
Got order:{CBG 0 1.05 100}
Got order:{CBG 1 1.05 200}
Matching:{CBG 0 1.05 100} against:CBG: 1655 0.8704 1.4305 2797
Matching:{CBG 1 1.05 200} against:CBG: 1655 0.8704 1.4305 2797
Matching:{CBG 0 1.05 100} against:CBG: 9414 0.8847 1.9453 9074
Matching:{CBG 1 1.05 200} against:CBG: 9414 0.8847 1.9453 9074
Matching:{CBG 0 1.05 100} against:CBG: 8787 0.9490 1.0373 12530 - Execute Order!
Matching:{CBG 1 1.05 200} against:CBG: 8787 0.9490 1.0373 12530
Matching:{CBG 1 1.05 200} against:CBG: 5840 0.9519 1.5722 12054
Matching:{CBG 1 1.05 200} against:CBG: 2485 0.9425 1.7220 11042
Matching:{CBG 1 1.05 200} against:CBG: 5235 1.0583 1.5840 4333 - Execute Order!
Labels:
channels,
golang,
goroutines,
pub/sub
Thursday, 5 September 2013
Sorting share trading books with Golang
Golang's sort package provides a set of primitives that allows us to define the order of our own structures. Furthermore we can define various ways of sorting the structures by extending the base collection.
type Order struct {
Ind string // buy sell indicator
Epic string
Price float64
Quantity float64
Time int64
}
type Orders []Order
func (o Orders) Len() int {
return len(o)
}
func (o Orders) Swap(i, j int) {
o[i], o[j] = o[j], o[i]
}
type BuyOrders struct {
Orders
}
type SellOrders struct {
Orders
}
func (o BuyOrders) Less(i, j int) bool {
// buy orders are reversed
return true
}
return false
}
// because of this element we cannot simply call reverse
return o.Orders[i].Time < o.Orders[j].Time
}
func (o SellOrders) Less(i, j int) bool {
return true
}
return false
}
return o.Orders[i].Time < o.Orders[j].Time
}
Labels:
golang,
share trading,
sorting
Tuesday, 14 May 2013
The Reliability of Go
As
part of the Canonical Cloud Sprint taking place in San Francisco last
week I attended Dave Cheney's talk at the GoSF meetup on the porting and
extension of juju. Juju is an open-source cloud management and
service orchestration tool that if you haven't heard of yet, you soon
will have.
After the talk an audience member asked if Go was reliable. Having used Go in production for coming up to three years now, without incident, this came as a bit of a surprise to me. Prior to moving to Canonical I worked for one of the UK's largest market makers. A market maker is basically a wholesaler for institutional share traders and stock brokers. During my time there I replaced several key systems components with Go.
System monitoring.
The services within the system were monitored by a python script, pinging each node, discovering services, connecting the networking dots, checking health etc. Due to the complex nature of the system this script could take up to three minutes to scan nodes and process the results. The script would often stall whilst processing the vast amounts of data produced. After porting the script to Go the runtime was reduced to under one second, and we never saw a single stall when processing.
Data store.
A legacy relational database was replaced with a Go based key/value store to remove bottlenecks at market open. This service is now the key piece of architecture in the system, processing all inbound and outbound quotes/orders to and from the London Stock Exchange, the Multi-lateral Trading Facilities, and key exchanges across Europe. This service processes instructions at an average of 7 microseconds (actually, 6 under Go1.1), and never once failed, even at peaks, processing tens of thousands of instructions per second. Go is currently providing key infrastructure components within the finance industry.
As I left my old position I was in the process of swapping the messaging middleware and the third-party price feeds with services written in Go.
After the talk an audience member asked if Go was reliable. Having used Go in production for coming up to three years now, without incident, this came as a bit of a surprise to me. Prior to moving to Canonical I worked for one of the UK's largest market makers. A market maker is basically a wholesaler for institutional share traders and stock brokers. During my time there I replaced several key systems components with Go.
System monitoring.
The services within the system were monitored by a python script, pinging each node, discovering services, connecting the networking dots, checking health etc. Due to the complex nature of the system this script could take up to three minutes to scan nodes and process the results. The script would often stall whilst processing the vast amounts of data produced. After porting the script to Go the runtime was reduced to under one second, and we never saw a single stall when processing.
Data store.
A legacy relational database was replaced with a Go based key/value store to remove bottlenecks at market open. This service is now the key piece of architecture in the system, processing all inbound and outbound quotes/orders to and from the London Stock Exchange, the Multi-lateral Trading Facilities, and key exchanges across Europe. This service processes instructions at an average of 7 microseconds (actually, 6 under Go1.1), and never once failed, even at peaks, processing tens of thousands of instructions per second. Go is currently providing key infrastructure components within the finance industry.
As I left my old position I was in the process of swapping the messaging middleware and the third-party price feeds with services written in Go.
Go's adoption is gathering pace thanks to the terse syntax, straightforward powerful
standard library, excellent tooling and concurrency primitives.
Go shows real maturity beyond its relatively young age due to the experience of the core development team and the consideration that is shown when introducing language constructs and extending the standard library.
I changed positions so that I could work with Go full-time. Ask anyone that knows me and they'll tell you that I'm not a betting man; you better believe Go is reliable.
Go shows real maturity beyond its relatively young age due to the experience of the core development team and the consideration that is shown when introducing language constructs and extending the standard library.
I changed positions so that I could work with Go full-time. Ask anyone that knows me and they'll tell you that I'm not a betting man; you better believe Go is reliable.
Labels:
go,
golang,
reliability
Saturday, 16 February 2013
Waiting for Golang channels to drain
Golang's channels easily map onto the producer consumer pattern. Lets assume that everything that is produced needs to be consumed, even if the process receives a SIGTERM.
The following example shows how we can register channels, monitor for a kill signal, and then wait for everything to be consumed.
package main
import (
"log"
"os"
"os/signal"
"reflect"
"syscall"
"time"
)
var (
BufferSize = 512
MaxIter = 10
monitored []interface{}
c = make(chan int, BufferSize)
stopping bool
)
func RegisterChannel(i interface{}) {
monitored = append(monitored, i)
}
func MonitorSigTerm() chan bool {
s := make(chan os.Signal, 1)
b := make(chan bool)
signal.Notify(s, syscall.SIGTERM)
go func(c chan os.Signal, b chan bool) {
_ = <-c
log.Println("Cleaning up")
// tell the caller
b <- true
for _, i := range monitored {
ch := reflect.ValueOf(i)
if ch.Kind() != reflect.Chan {
continue
}
prev := 0
iteration := 0
for {
if ch.Len() == 0 {
break
}
if prev == ch.Len() {
iteration++
// enough?
if iteration >= MaxIter {
log.Println("Dropping")
break
}
} else {
iteration = 0
}
prev = ch.Len()
log.Printf("Draining:%v\n", prev)
// other goroutines are working, let them
time.Sleep(1e9)
}
}
os.Exit(1)
}(s, b)
return b
}
func main() {
RegisterChannel(c)
stop := MonitorSigTerm()
go func() {
i := 0
for {
if stopping {
break
}
i++
c <- i
time.Sleep(1e9)
}
}()
go func() {
for {
i := <-c
log.Printf("rx:%v\n", i)
// slower read
time.Sleep(2e9)
}
}()
stopping = <-stop
// wait for cleanup to finish
select {}
}
The following example shows how we can register channels, monitor for a kill signal, and then wait for everything to be consumed.
package main
import (
"log"
"os"
"os/signal"
"reflect"
"syscall"
"time"
)
var (
BufferSize = 512
MaxIter = 10
monitored []interface{}
c = make(chan int, BufferSize)
stopping bool
)
func RegisterChannel(i interface{}) {
monitored = append(monitored, i)
}
func MonitorSigTerm() chan bool {
s := make(chan os.Signal, 1)
b := make(chan bool)
signal.Notify(s, syscall.SIGTERM)
go func(c chan os.Signal, b chan bool) {
_ = <-c
log.Println("Cleaning up")
// tell the caller
b <- true
for _, i := range monitored {
ch := reflect.ValueOf(i)
if ch.Kind() != reflect.Chan {
continue
}
prev := 0
iteration := 0
for {
if ch.Len() == 0 {
break
}
if prev == ch.Len() {
iteration++
// enough?
if iteration >= MaxIter {
log.Println("Dropping")
break
}
} else {
iteration = 0
}
prev = ch.Len()
log.Printf("Draining:%v\n", prev)
// other goroutines are working, let them
time.Sleep(1e9)
}
}
os.Exit(1)
}(s, b)
return b
}
func main() {
RegisterChannel(c)
stop := MonitorSigTerm()
go func() {
i := 0
for {
if stopping {
break
}
i++
c <- i
time.Sleep(1e9)
}
}()
go func() {
for {
i := <-c
log.Printf("rx:%v\n", i)
// slower read
time.Sleep(2e9)
}
}()
stopping = <-stop
// wait for cleanup to finish
select {}
}
Labels:
channels,
golang,
reflection,
signals
Tuesday, 29 May 2012
Straightforward, Powerful, Standard Library
Whilst testing a Go application to destruction I was manufacturing various scenarios and needed to discover where exactly each goroutine was in its processing. After a quick scan of the documentation, and 5 minutes keyboard bashing I was able to add a couple of simple functions to a utility package and enable applications to dump their current stack when signalled.
func ListenForPrintStackSignal() {
sigCh := make(chan os.Signal, 1)
signal.Notify(sigCh,syscall.SIGUSR1)
go func(sigCh chan os.Signal) {
for _ = range sigCh {
log.Print("Request to print stack")
PrintStack(true)
}
}(sigCh)
}
func PrintStack(all bool) {
//arb sized buffer
buffer := make([]byte, 10240)
numB := runtime.Stack(buffer,all)
log.Print(fmt.Sprintf("%v current goroutines",runtime.NumGoroutine()))
log.Print(string(buffer[0:numB]))
}
Go's standard library is by far the most powerful and striaghtforward set of builtins I have used. Everything seems to be in place (or planned for) to quickly plug functionality together without having to jump through hoops or battle oddly named methods with bizare signatures. The productivity kick that comes with using this standard library is one of the reasons that we will see an exponential uptake of the Go language in the coming months and years.
Thanks to @rogpeppe for the golfing tips.
Labels:
go,
golang,
signals,
standard library
Tuesday, 22 May 2012
Simple Synchronous and Asynchronous Logging
Almost all services need to write to a log for auditing purposes. When these services require high throughput this write to disk can be a real bottleneck.
So, I wrote a simple log package that implements asynchronous logging to standard output. I only ever write to standard out for service logs. For me redirecting to a file from the command line has proven too flexible over the years to even consider named logs.
Here is the ad/log package:
/*
Package log implements standard logging functions for logging to standard out
*/
package log
import (
"log"
"os"
"fmt"
"strings"
"time"
)
var myLog *log.Logger = log.New(os.Stdout, "", log.Ldate+log.Lmicroseconds)
var (
AsyncBuffer = 1000
asyncLogging = false
asyncChannel chan string
)
//Print message to standard out prefixed with date and time
func Print(s string) {
if asyncLogging {
s = "(" + strings.Fields(time.Now().Format(time.StampMicro))[2] + ") " + s
asyncChannel <- s
} else {
myLog.Print(s)
}
}
//Enable / disable asynchronous logging
func EnableAsync(enable bool) {
if enable == asyncLogging {
return
}
if enable {
asyncChannel = make(chan string,AsyncBuffer)
go func() {
for {
message, ok := <-asyncChannel
if ok {
myLog.Print(message)
} else {
asyncLogging = false
return
}
}
}()
} else {
close(asyncChannel)
}
asyncLogging = enable
}
//Returns current asynchronous logging value
func Async() (bool) {
return asyncLogging
}
Here is a simple example:
package main
import (
"fmt"
"time"
"ad/log"
"ad/utl"
)
var i = 0;
func main() {
utl.UseAllCPU()
incAndPrint()
log.EnableAsync(true)
for {
incAndPrint()
incAndPrint()
fmt.Printf("Doing my own thing\n")
incAndPrint()
time.Sleep(1e9)
}
}
func incAndPrint() {
i += 1
log.Print(fmt.Sprintf("%v",i))
}
Output looks like this:
~/go/src/ad/misc$ go run logT.go
2012/05/22 15:39:26.229273 Running log 1.2.1
2012/05/22 15:39:26.229771 Running utl 1.0.0
2012/05/22 15:39:26.229812 Setting UseCPU from 1 to 4
2012/05/22 15:39:26.229874 1
Doing my own thing
2012/05/22 15:39:26.230098 (15:39:26.229989) 2
2012/05/22 15:39:26.230151 (15:39:26.230035) 3
2012/05/22 15:39:26.230163 (15:39:26.230092) 4
Doing my own thing
2012/05/22 15:39:27.230701 (15:39:27.230555) 5
2012/05/22 15:39:27.230761 (15:39:27.230612) 6
2012/05/22 15:39:27.230775 (15:39:27.230696) 7
Doing my own thing
2012/05/22 15:39:28.231295 (15:39:28.231143) 8
2012/05/22 15:39:28.231356 (15:39:28.231203) 9
2012/05/22 15:39:28.231369 (15:39:28.231262) 10
For me the power of this painfully simple package comes when we code log.Print and take the value of the log.asyncLogging flag from a command line switch. This allows you to code in a standard fashion but run individual service instances up with a dedicated logging goroutine if entering a busy period. Alternatively, the log.asyncLogging flag could be set on the fly when the service decided; say when demand is high.
Saturday, 12 May 2012
System monitoring
I hate having to configure monitoring. Especially in a complex interconnected system with 100s of servers and 1000s of processes. I have used Nagios in the past; this might be why I'm bald.
When monitoring a system all I really want to know is what should be up, what is up, what is it connected to, what are the inputs / outputs, and how busy is it.
A while ago I wrote a Python script to take a list of hosts, ssh to them, issue a ps and lsof, parse the outputs and connect all the dots. This worked well, but it took a long time to run and output a massive file giving a lot of data (network connections, stdin/out, executables, command line arguments, cpu / memory usage) but no real useful information. Apart from lengthy execution and no information it worked a treat. So, having failed to acheive anything useful I naturally put this python code into my back-burner directory and got on with something else.
A few months ago I read somewhere that the Processing environment was now available in the broswer via processing.js. This was good news. I now had the tool to allow me to successfully display the data as meaningful information. I just now had to get the get the server scan down to a sensible time in order that the information displayed was up to date. I also noticed that using python to monitor my system's performance was causing systems performance issues!
Enter Go.
I ported the python code to Go. Not a straight port; I utilised goroutines to issue the ssh commands and channels to gather / process the output and push the information to browser clients.
On a loop, ssh commands are issued and the output processed into meaningful results which are then stored in a map. Process signatures are parsed so as to give sensible names and display arguments. This map is then compared to the previous loop's map to calculate each process's status. Is the process alive? Has the PID changed? Have the socket connection changed, or dropped? How much CPU and memory are being used. That sort of thing. The results are then pushed as JSON to any connected browsers via websockets.
In the browser, the processing.js environment loops at a given frame rate drawing the nodes and connections as detailed in the provided JSON. The screen gives lots of visual information: status, type of process, resource usage, connections made etc. The Processing and jQuery allowed me to code user interactions with the nodes. The user can position, filter, detail, and maintain the nodes. As the websocket is bidirectional and a node details the process's standard output I was able to allow the users to hit a key to tail the process's log file in the browser. I intend to expand this further to interact with the underlying process; start stop etc, but I haven't added any authentication or authorisation yet.
Once the user has positioned the nodes to their liking they're able to save the screen by sending the current display context, again as JSON, back to the server over the websocket. This then allows me to save a good known state of the system and use this as a default process map when starting the monitoring process. This solves the problem of monitoring missing processes.
Processing enabled me to present the system information in a palatable way and deserves a lot of credit.
Go though, deserves unending credit. All the heavy lifting of issuing 1000s of system calls, parsing the output and collating the results is handled with ease. The job that took my python script four to five minutes to complete is done in under a second. I even restricted the Go process to a single core to make the comparison with python fair. The Go garbage collector is clearly doing its job as throughout the day the memory footprint holds steady. Sar stats across all machines show no adverse effect of using this scatter / gather approach.Go, Processing (along with jQuery for the dialog inputs) and the good old shell utils gave me the ability to write my own bespoke monitoring software and ditch Nagios.
My one complaint though is that no matter how much I rub the source code on my head my hair still wont grow back.
Labels:
go,
golang,
nagios,
processing,
processing.js,
system monitoring
Subscribe to:
Posts (Atom)