Can Spark Streaming Survive Chaos Monkey?
techblog.netflix.com
techblog.netflix.com
I've had problems with executors dying when running under YARN, and it makes keeping track of running jobs and finding log output more difficult. So unless you really need to run on the same instances as other MR tech, it seems standalone is the way to go.
If only AWS provided EMR AMIs for spark standalone clusters I could switch...
The problem I found seemed to have been discovered by one or two other people on mailing lists - it seemed something to do with YARN terminating spark executors for memory usage. It was odd since I wasn't using the cluster for anything else - YARN was only used as the scheduler because it was installed in the AMI, and I didn't want to use the ec2 scripts to start a standalone cluster just on EC2s.
It was difficult to debug because YARN makes things pretty opaque when tasks just die - logs weren't copied around the cluster so I had to SSH to the instance that died and try to find something to indicate what had happened, and half the time the error messages were vague.
I also didn't work out how to view the job tracker UI when running under YARN. "Learning Spark" says I need to proxy through the YARN cluster somehow but doesn't give an example... So I have no progress about running jobs at the moment, only about finished jobs.
In the end, I upped the instance size I was using and everything went well. I think my dataset must have been too large and spark was spilling to disk which tripped things up. So at the moment I don't have confidence that I could process huge datasets with Spark on EMR which is a bit annoying :-(
- set up N available servers
- make clients store to N servers at the same time when they calculate the value
- query M (where M<=N) servers before deciding you have to recalculate
- If you got a response from server between 2 and M, re-store the value everywhere (just pay attention to preserving timeouts)
And you get distributed, self-healing, chaos monkey resistant memcached without any support on the server side.
Also if you want to avoid stampede, you could insert 1s TTL placeholders that mean "back off, someone else is calculating" into keys you know are popular and may experience contention. Just make sure you use CAS so you don't overwrite data with placeholder.
(I worked on the system)
There is also another project called Dynamite [0] that puts a gossip/Cassandra-like protocol in front of redis.
[0] http://techblog.netflix.com/2014/11/introducing-dynomite.htm...
I spent the first two weeks clarifying that I was indeed referring to Apache Spark, not the Spark-Formerly-Known-As-Telecom. They'll always be Telecom to me.