Twitter open-sources a high-performance replicated log service
github.com
github.com
I think the motivations they list are: 1. Different I/O model 2. Was started before Kafka had replication (the first release of Kafka with replication was in late 2013 I think)
The I/O model I'm less sure about, we looked at similar things for Kafka and they didn't seem worth it (basically you're doing a ton of stuff at the app level that the OS does pretty well--namely caching and buffering linear I/O), we'd have to look at actual benchmarks to know.
Here is my take on the pros and cons of the core tech.
Pros: - Seems to have better built in support for fencing/idempotence - Better geo placement?
Cons: - Lots more moving pieces. Already people are irritated that there are both Kafka nodes and ZK to set up. This system seems to split this over separate physical tiers for serving, core, storage, and zookeeper. My experience has been lot's of tiers is generally a big headache.
Neutral: - There seems to be a built in achival to HDFS. I think if the consumer is fast and efficient then you don't need to reach around your consumer api which will be high latency (since you have to wait for files to be closed out).
There is also a bunch of stuff Kafka does that I'm just not sure about how complete it is in DistributedLog: - Clients in a bunch of languages - Integration with all the major stream processing frameworks - Log compaction http://kafka.apache.org/documentation.html#compaction - Connector management http://www.confluent.io/blog/announcing-kafka-connect-buildi... - Quotas/throttling - Security/ACLs
Also, once you're on the Big Data™ train, a lot of things like to plug into Zookeeper, so it becomes more of a convenience.
Kafka, and presumably DL, are at their most useful when you're pushing the limits of NIC and/or HDD performance for throughput. Zookeeper's configuration is a footnote in the complexity of managing one of these systems, and lets them avoid implementing their own byzantine coordination system. Also, folks seem to appreciate Aphyr's opinion, and he states it pretty plainly: Use Zookeeper. It’s mature, well-designed, and battle-tested. [2]
* attempting thousands (tens of thousands?) of simultaneous writes across data centres/continents at the same instant and then saying it was slow. Subscribing all clients to updates on all parts of the tree. The architecture had been grown by people who didn't understand the guarantees and constraints.
* calling it unreliable after arbitrarily moving nodes around without changing connect strings and generally mis-configuring it. Essentially, pointing clients at machines that no longer contained nodes and blaming ZK for this not working.
It could be easier to setup. People often don't want to think about such things at all. But it's also not the hardest, and I found it to be very resilient to node and network failures when deployed correctly.
Which is just to say there's room for improvement. A dependency on Zookeeper is fine if you've already got a configured cluster, and a cognitive speed bump if not.
On another note, I find it somewhat funny that these are called "log" services, logging is probably the least interesting use case for these things I can think of. A better description in my mind would be as a distributed event processing framework, since what they are really doing is distributing discrete events in a reliable manner.
https://en.wikipedia.org/wiki/Log-structured_file_system
The use of the term "log" in some contexts implies concepts like lossy, unimportant, non-durable, etc. However, in the context of log-structured file systems or journaling file systems, the sense that's intended is append-only as a primitive, atomic, consistent, durable (replicated) operation.
https://engineering.linkedin.com/distributed-systems/log-wha...
https://engineering.linkedin.com/distributed-systems/log-wha...
An ordered append only datastructure is rightfully called a log. The fact that text based files are called logs is just an annoying feature in common english usage.
Kafka has the open-source/java community and ecosystem that fits in well with the rest of the current big data processing stuff though.
What's your unique ID scheme?
Let's say I'm willing to believe[1] that you've got Durable and Consistent down, once messages make it committed in to the system. What's the story for messages on their way in? My application logs are buffered to the local disk, now I'm streaming them into central storage, and halfway through a TCP connection that's shuffled 2mb of thousands of messages into storage, the connection terminates -- unexpectedly, midmessage. Could the service have committed more messages than it acknowledged? Or many less than I've sent? Both could be true from the network standpoint.[2]
So, what I need to know, and what should be very easy to answer, front-and-center in your docs, pretty please:
1) Where should my log uploader resume?
2) Is there any danger of repeatedly entering some lines?
3) If I have log lines that are legitimately duplicates, will they be stored at the correct count?
These are questions that may have a different answer than the durability after data makes it fully into the system. It also may provide useful information about how complexity the code in a submitting client is, because good answers tend to require some kind of ID sequence being assigned on submitting clients, afaict. And it's really just plain critical to sanity.
----
[1] well, no, I'm not, "trust by verify" in all things etc etc; but let's suppose that's more believable and something I have to mechanically verify anyway, and doesn't have an obviously observable boolean at the protocol level as to whether it's going to work well or not, and system internals simply don't have such a sordid history of being over-simplified until they're broken like client interfaces so often are, so...! We'll handwave that to a later and more involved step of quality investigation.
It looks like they use fencing and a two-phase commit to prevent duplicate writes. Whether that covers all failure scenarios I'm not sure.
Yes of course. It may have committed records but not had time to to ACK success.
> Or many less than I've sent?
Yes again. However it will NEVER ACK success for writes which it hasn't committed.
So the persistence story is pretty straightforward (I'll see if we can update docs if this is missing).
If you want to avoid duplication on the write side, you would have to retrieve the last log record-- with id x-- and start adding records again from id x+1.
Fencing + locking combine to provide efficient exclusive access and the fat client (for now) guarantees write ordering.
I have had too many deployment nightmares with Zookeeper. I would prefer to avoid it as much as possible, plus systems software in Java, sigh.
So running zookeeper in a dynamic environment becomes very risky. Aka AWS. Perhaps on GCE it's less dangerous, because the host migration is pretty good.
There are other operational issues people have noted above as well.
You also need to disable the JVM's insane DNS-caching behavior (config change only).
Lots of things can and do go wrong with Zookeeper. I suspect it depends on the use case, but building a Zookeeper dependency into any system is potentially asking a lot of users/operators.
Ideally, these kinds systems should be built such that the distributed coordination piece is pluggable and the implementation can be chosen based on deployment concerns.
On the other hand, as a mortal, I can write a-grade-above-code-that-an-idiot-would-write-just code in Java at about 10 times the speed I can write dire-useless-risible C code.
The comparison is pointless though, good modern languages like Rust and Julia are developing and LLVM is enabling further development.
Also even IF you are tuned properly, there is a world of 99-percentile you'll never get to.
Additionally, the difference between openJDK and Oracle JDK becomes something admins learn to hate you for. Complex shell scripts to invoke java, and lots of standard unix tools just don't work super well. eg: pgrep and pkill. you can tweak it, but it takes a little while to learn the many tips and tricks.
It's not all bad, most Java programs are deployed in a static-all-batteries-included fashion, so you rarely worry about system-installed library versions. So that makes deployment a little less hassle. You never have to recompile for cross platform. The profiling tooling and other stuff is pretty good, and the more you're willing to pay the better the tools get.
The point is you cant use binary name to make your way around anymore. The binary name is 'java'. Standard unix tools just don't work. Yes there are work arounds, but over time things end up being just a little more complex than they should be.
Which means if someone is comparing zookeeper and etcd, well etcd wins major points for being unixy and easy to deploy. Copy 1 binary, done. ZK loses major points here. Gotta make sure the JVM is installed, but do you need the OpenJDK or the Oracle one? if the latter, well apt-get and yum are less helpful.
It's all just little globs of annoying details that add up to be a small pain. Nothing horrible, but if you could make a choice to avoid that, why not?
Basically I guess what I'm saying is there is probably a market for replacing all the Javay distsys stuff with Rust/Go versions. I mean look at etcd!
Or in the other direction. If your main system runs on the JVM and your sysadmins are used to the JVM tools then having a piece of infrastructure that's just another .jar is wonderful, and C/Rust/Go/Ruby/etc. infrastructure elicits groans. Mixing platforms will always be harder than a common platform. So the infrastructure market depends on where you think the future of applications is.
In terms of simplifying the deployment, they could have went embedded with Atomix[1] instead. Perhaps next time.
[1]: http://atomix.io/
When I saw the header I thought to myself, "Yay! Now I can run something like Kafka, without depending on Zookeeper or having to install JVM!! Give me stable drivers for Python and Go and I am sold!"
If I have to install Zookeeper and JVM, why not use Kafka?
Seems reasonable, right?
Except "[3] Kafka addressed these durability concerns in version 0.8"
So they built a whole thing, because they didn't bother to ask or say "hey, if we help fix the durability, would that we welcome?" or even "do you guys have a plan and timeline to fix it?"
That's .... not great.
Now, maybe there are other reasons, but they aren't elucidated in this blog post :P.
Even then, my general view would be "did you approach the community and discuss your concerns or just dismiss them out of hand as infeasible", mainly because my experience is that if you do this, you often find they have exactly the same set of concerns/goals, and just need more resources to make it happen.
The desire to build shiny objects is very large. Outside of the paper plans of engineering teams, these things rarely end up up more shiny than what already exists or will be built by the time you are done.
The software industry has a real problem contributing to open source, the stuff, you know, allowing them to make money in the first place.
What's the point of retaining an engineer who's doing nothing for the business but re-inventing existing successful software?
actually, trying to let people do work you don't need done, in order to keep them happy, is a pretty rookie manager mistake.
If they aren't passionate about it, and you can't persuade them to do the things you need doing, they aren't the right person for the job.
That is always true, even if they were the right person in the past.
Your goal in that case should be to try find stuff the company needs done that they want to do, and push them to work on that. But if you find nothing, ...
Not saying that they couldn't have worked more closely with Kafka's team in the first place, but, hey, now we have two Kafkaesque log services instead of just one. Seems like a win to me.
There are also drawbacks to consider: Those experts might just decide to leave and then you have an in-house solution that you have basically no chance to find experts on. At least with an open source solution you may have an easier time.
CantTellIfSerious.jpg
Since I'm on AWS EC2, I want to try this:
- Write the logs to local SSD, asynchronously
so as not hold back the http request.
- Have a separate cron job that loops through
the log directory and scoops up all the files.
- The job will then stuff those files into a Kinesis Firehose.
AFAIK, Kinesis Firehose does not require any capacity provisioning,
unlike the Kinesis Streams, so I'm set "for life" (up to 5MB/second)
- The firehose will accumulate the logs and put them into S3.
Hurray unlimited storage!
- S3 will trigger a Lambda.
- Lambda will parse through the log from S3, pull out
interesting properties (IP address, user id, session id, etc) and
stuff them into a DynamoDb table.
- If I need to see data from one user/ip/session I will use DynamoDb
to find the right S3 blobs.
- If I need to reprocess the logs to extract a new piece
of data that I did not foresee earlier, I can run a
map-reduce task
Except the last piece, this looks like something I can half-ass in a couple of days and forget about it for another couple of years.Any opinions? I don't really want to use a SaaS log service because gigabytes per day.
I've used Lambda a bit. The debugging process can be a pain, since you're forced the upload a ZIP file, and if your code times out Lambda doesn't give you any traceback to indicate what happened. There's also a maximum run time of each Lambda invocation, which I believe is 5 minutes. Is there a chance your parsing may run longer than that? Also, what will you do if you upload some bad code and the parsing fails? Will it be the end of the world if you lose data while you fix the parsing?
Oh, I see you plan on doing map-reduce to re-parse the logs, so maybe that part isn't as big a deal.
You could also consider doing something like rsyslog -> db-of-choice while also rotating the files off to S3 for long-term map-reduces. This is all to ignore the obvious ELK cluster solution, which will give you good data visualization and investigation options, but may be more of a headache to set up and maintain than you are looking for.
Anyway, those are my thoughts. Hope they help.
For Kinesis, I planned to use Firehose, not Streams (the latter have to be provisioned, which I was hoping to avoid). The firehose could put data into S3 for me, and S3 would trigger lambda. However I just realized that S3 will only make 3 attempts to invoke the Lambda, so that pretty much rules out this part of my design - the data will not get lost, but it will not get indexed either. I may run map/reduce later, but I don't want to be dependent on doing that to pick the loose ends.
These are just server logs, they don't affect business continuity. Still I wouldn't want to be sitting there and wondering "is this user having connectivity problems, or did I just lose a pile of logs?".
I could probably put the files directly into S3, got carried away stacking my AWS features together. :) I'd need to be more careful with batching, so as not to create batches too small or too large. Perfectly doable, though Kinesis Firehose already does that for me. Plus in case of the EC2 instance death I will lose the current batch with the hand-rolled solution, but not with Firehose.
So I guess I should just put a batch of data directly to S3 and send an SQS message to make sure that it's indexed properly, then delete the local files.
I really don't want to pick up another thing that I have to understand and manage. Like ELK. Someone who knows ELK will probably have no problem managing it, but my head is full with business domain problems.
// Create Stream `basic-stream-{3-7}`
// dlog tool create -u ${distributedlog-uri} -r ${stream-prefix} -e ${stream-regex}
./distributedlog-core/bin/dlog tool create -u distributedlog://127.0.0.1:7000/messaging/distributedlog -r basic-stream- -e 3-7
You use a regex to progammatically create N streams, which would be N shards in Kinesis or N partitions in Kafka.Depending on use-case, of course, there's lot of places where scala is preferable to Java. A byte-shuttling service where performance is important, mutability reigns and expressiveness is irrelevant.. not one of those places.
Joke aside, at least in hiring people it should.
Cassandra doesn't work, and Hadoop is a complete waste of hosts for most companies (hence the move to Spark.)
...which also runs on the Java Virtual Machine and is subject to the same pros and cons.
Hadoop is slow, but on huge data volume the overheads are dwarfed by the parallelism gained. Most companies don't have huge volume though.
For example recently I saw someone propose using Hadoop for a sub-TB dataset...
So, what data do you base these assertions on ? Also, not to burst your bubble but a lot of businesses (if not the majority) run Spark on YARN. And Spark is built on the JVM.
We detached this subthread from https://news.ycombinator.com/item?id=11669189 and marked it off-topic.