15 ms·
Publishing with Apache Kafka at The New York Times
- qaq 9y agoResume driven architecture taken to extreme.
- manigandham 9y agoThe very definition of over-engineered. This is just event-sourcing turned into a marketing article for Kafka. It doesn't really matter what the "source of truth" system is, although Kafka doesn't seem quite as mature/stable enough for that. With such little data, they can push into a nice graph database instead, run it entirely in memory and meet 10x any demand they'll ever see along with all the query types they'll need. Add in an elasticsearch cluster on the side and problem solved. Any database can serve as a log replay, as long as you save all the versions then it's just called a query.
- robbyt 9y agoWhy do you think it's over engineered? Don't you think the Times have considerable capacity requirements?
- manigandham 9y agoThey claim in the article to have 100GB of text data. Let's bump that up 2 magnitudes for all text and metadata ever (outside of media files) and you can still run the entire thing on a single rack of servers and meet any performance needs. Many industries and applications are leagues ahead in both data size and speed - this isn't an example of such.
- seanp2k2 9y ago100GB...let me go find the largest MicroSD card in my house to put that on. That would actually be a good fit, since it'd work in an rPi3 which could likely serve their data publishing needs (assuming only a few updates to articles per second..not mentioned in the article, but I'd be surprised if there's more than that given the data sources) vs what they've done. Honestly, what kind of RPS are they talking here? Requests per minute, if that, seems like.
- John23832 9y ago>The very definition of over-engineered. I agree, it is a lot. However, it is an interesting approach to their problem. The over engineering makes the flow much easier to handle. >Any database can serve as a log replay, as long as you save all the versions then it's just called a query. The article addressed their potential issues with snapshotting
- manigandham 9y ago> The article addressed their potential issues with snapshotting They say it would be outdated - but if you store all the versions like I said, what is getting outdated? Kafka consumers are essentially doing the same thing - it's a poll-based model that asks for more data from the current offset, no different than a SQL query with a where clause.
- vikiomega9 9y agoUnless I'm mistaken, If I were to build out a simple event log represented by a relational DB, I have bottle necks when writing to it, and have lag in terms of processing the events, and if I were also pushing those events to a queue to hydrate aggregate snapshots I would have to have client logic to deal with duplicate events or not acking processed events etc? Intuitively, I guess kafka is "more realtime" and "more available" when compared to the home-brew event log? EDIT: obviously those constraints in my home-brew event log are relaxed when my problem domain is amenable to things like associative operators, idempotency, inverses etc.
- cookiecaper 9y agoThere's no reason not to make good use of Kafka or similar solutions. The issue is that people use it without understanding it. In this article, they say that Kafka is their system of record and their primary long-term storage. That's very silly.
- subsubsub 9y agoWhy is it silly to use as system of record or long term storage?
- qaq 9y agoit's exact opposite the main cost is gc of dead rows in MVCC RDBMS since you are never deleting or updating rows performance will be very decent for writes.
- manigandham 9y agoYou can use Kafka as the buffer/processing log before persisting to the database, but with such a small dataset it's just not necessary. It's a news publishing system, not high-frequency trading.
- vikiomega9 9y agoWell my point is that it's probably faster to get to production if I simply used Kafka _when modelling my work flow as an event processing system_ but it took them a year so I don't know now haha
- ma2rten 9y agoYou could use any database, but Kafka has apis specifically for that. Why would you reinvent the wheel? What makes you think that Kafka is unstable?
- manigandham 9y agoEvent sourcing is unnecessary here. What they want is a versioned history of their content with flexible schemas and arbitrary queries. Instead of using a strong fast graph database as a perfect fit, they chose to implement it poorly using Kafka which I see as much more time wasted on reinventing the wheel. The dataset is small and all of the "very different" use-cases are just downstream apps that query a database. Why use Kafka to then materialize several different databases when a single graph database can serve all these downstream apps? Remember Kafka consumers themselves are just polling queries against a log. A processing log is one thing but event-sourced source of truth in custom database logic in Kafka is borderline ridiculous. This project is more moving parts and less functionality while cleaning none of the existing mess. Effort would've been better spent consolidating all their systems instead (which they still have to do since these standard schemas need to be used somehow).
- puranjay 9y agoReminds me of the NYT "Snowfall" story website that NYT spent $100k+ on producing, and some startup remade within a couple of days on Wordpress
- josephg 9y agoThat doesn't prove that the NYT's money was poorly spent. Decisions are always slow and expensive, and typing is quick. Copying an existing design is almost always cheaper than making the original. And in many cases, copying is dramatically cheaper. My favorite example is the iOS game Threes[1], which took 14 months to make and then was cloned on android in ~20 days. It was cloned so quickly that the original developers were accused of copying the cloned version! But just because the cloners made a functionally identical product doesn't mean they did the same work as the original designers. They got to skip designing anything - which is usually the hardest and slowest part. [1] http://asherv.com/threes/threemails/ http://asherv.com/threes/threemails/
- neeleshs 9y agoI was about to comment that while any database can serve as a log replay and may even be feasible for these volumes, scaling it at high volume later would be extremely hard - I have experience with this in a multi-petabyte set up and its a nightmare using an RDBMS. But, the article says "In Apache Kafka, the Monolog is implemented as a single-partition topic" - losing all the goodness Kafka provides around scaling. This setup now is no better than an RDBMS with master/slave replication.
- mehh 9y agoWhy not just have the 'logs' in s3 in sorted/indexed buckets.
- LaFolle 9y agoIf all articles are put into monolog, what is the procedure to fetch all articles, lets say, published in year 1857? Will that be a O(n) operation (assuming all messages published in monolog have a timestamp field).
- khlbrg 9y agoI also recommend the talk from Kafka summit. https://kafka-summit.org/sessions/source-truth-new-york-times-stores-every-piece-content-ever-published-kafka/ https://kafka-summit.org/sessions/source-truth-new-york-time...
- deleted 9y ago[deleted]
- oliveralbertini 9y agoInteresting, we use rabbitmq instead of kafka and we have a re-indexation system... not sure if it's more complex for what I see.
- pm90 9y agoExcellent, well written article. The key take away seems to be that instead of an temporary event stream log, since the number of news articles (and associated assets) is finite and cannot explode, they store all the "logs" forever (I'm using the term log as is defined in the article, as a unit of a time-ordered data structure). I wonder if NYT can help other news websites by making their code open source? I'm a huge fan of NYT and their jump to digital has just been amazing. However, I would also like my local newspaper (which covers more regional news) to be able to serve quality digital content.
- dangayle 9y agoOne of the biggest issues is that a lot of newsrooms have zero control over their CMS, because they're owned by a corporate entity that dictates IT decisions from afar, slowly and with much gnashing of teeth. Family-owned papers like the one I work at (The Spokesman-Review in Spokane, WA) are the one of the few news orgs that could actually put something like this into production within this century, but even we still have to deal with manpower issues.
- srigi 9y agoWhat would be a bigger thing is if they opened access to their monolog.
- thinkMOAR 9y agoIt would make me very sad if this would be the case, that every newspaper is going to be running the same kind of software stack. It's then 1 step away from just having only 1 news paper in total. Diversity, choice, innovation, etc up to different reporters reporting on the same story with their point of view, will all disappear if all newspapers are going to run the same software stack.
- sync 9y agoNYT has quite a lot of open source repos: https://github.com/nytimes https://github.com/nytimes
- knowtheory 9y ago> I wonder if NYT can help other news websites by making their code open source? Hey! I, and a number of other news nerds have been encouraging FOSS for the past decade or so. And in fact a number of major open source projects have come out of news related projects, including Django, Backbone.js/Underscore.js, Rich Harris's work on Svelt.js, and a whole lot more. Most often the problem with local news organizations are operational constraints. The news biz has seen a huge downturn over this same period of time. Most orgs, both on the reporting side and on the tech side are super strapped for people-time. It's not enough to have FOSS software, you also have to have folks doing devops and maintaining systems often at below-market salaries.
- cookiecaper 9y ago>We need the log to retain all events forever, otherwise it is not possible to recreate a data store from scratch. SIGH. Cue the facepalm, head in hands, etc. I'm not going to get into a big thing here. But if you find yourself saying "I need to keep this thing forever no matter what" and then you try to use something that even entertains the notion of automatic eviction/deletion semantics as the system of record, you're doing it wrong. Not to burst the bubble of the techno-hipsters, but Kafka is "durable" relative to message brokers like RabbitMQ, not relative to a system actually designed to store decades of mission-critical data. Those systems are called "RDBMS". Elsewhere in the article he says that they have less than 100GB of data and that it's mostly text. This is massive overarchitecture that isn't even covering the basic flanks that it thinks it is, such as data permanence. I would really like to read the article that discusses why Postgres or MySQL couldn't have served this purpose equally well.
- bognition 9y agoOOC what makes a RDBMS more durable that a Kafka? Both of them are systems for representing data on disk. I'd love to hear why one representation system is better at disaster recovery than another.
- cachemiss 9y agoThat's a bit of an oversimplification. Production grade RDBMS systems have far more guard rails, testing and work put in to them than Kafka. It's relatively straight forward to lose data in Kafka, I've done it (its usually control plane bugs, not data plane).
- tepe1 9y agoGP is clueless. Kafka is in fact more resilient and much more performant than systems like MySQL. It's the difference between having one writer and many writers. More importantly, adding yet another database to an enterprise that's already overrun databases is not going to fix anything. Systems like Kafka are deployed precisely because there are already many domain-specific databases at work and now there's a need to tie them together. You can't have a single database that works for all domains, for everything from elastic search text to billing. That's just dumb. So you have many databases and tie them together using a transactional event log. BTW, the NYTimes, like virtually any mature enterprise, probably has a robust data warehouse and data retention strategy. I kind of doubt they're in the habit of losing data given their archives go back more than a century. But it's important to distinguish front-office transactional concerns (actually handing real-time requests from content producers and consumers) from back-office concerns (reporting, back-ups). They are very different domains as many enterprises are slowly discovering.
- tabeth 9y agoI wonder how much of this kind of stuff exists out of necessity and how much of it exists because very smart people are just bored and/or unsatisfied. Are there any articles that supplement this that explain how much business value is added/lost by the existence/removal of these kind of features? In the case of NYT I suspect its popularity is maintained because of the perception (real or not) of high quality journalism, in spite of any technical failings. --- How much would be lost if NYT was just implemented as text articles that are cached and styled with some CSS. "Personalization" could be added by tags each article has and a small component that shows the three most recent articles that share the same tag.
- weego 9y agoI can't speak for the situation at the NYT but the actual public site for online papers are often pretty simple things with most of the complexity being ad logic. The systems here almost entirely deal with writing and content retrieval pipelines for stuff that was written years ago in other systems that isn't tagged or stored in sympathetic ways, and there will also be the very old school print pipeline to have to deal with too.
- manigandham 9y agoHow much time/effort does it take to just update all this old stuff once and for all?
- dangayle 9y agoLiterally teams of interns manually re-typing old articles from microfilm. OCR isn't quite there yet, not for dealing with the ways newspapers and newspaper design has changed over the years. You'd get the text, but there would be no guarantees that it that story bylines, headlines, factboxes, and photos match. Our own digital archives go back to 1994, anything before that is manual input.
- cookiecaper 9y ago>I wonder how much of this kind of stuff exists out of necessity and how much of it exists because very smart people are just bored and/or unsatisfied. That's a ton of it. Like it or not, publishing a digital newspaper is not a hard or unsolved problem; it's one of the web's core competencies. If you hire people who want to build cool stuff to supervise a CMS, well, you get this kind of outcome. The raw cost is understated because these experimental setups misinterpret the functionality of the new architectures/formats they're using. It doesn't truly rear its ugly head until there is a major data loss or corruption event. It's not that these never happen with RDBMS, it's just that RDBMS contemplates this possibility and tries to make it pretty hard to do that, whereas message queues just automatically delete stuff (by design, so they can serve as functional message queues!). RDBMS have spoiled us and we take its featureset, 40+ years in the making, for granted. We need to be careful and not assume that `GROUP BY` is the only thing we leave on the table when we "adopt" (more accurately abuse) one of these new-wave solutions as a system of record. Since no one is going to admit to their boss "this wouldn't have happened if we used Postgres", and since most bosses are not going to know what that means, most of these spectacular failures will never be accurately attributed to their true cause: developers putting their interest in trying new things above their duty to ensure their employer's systems are reliable, stable, and resilient.
- jdcarter 9y agoFWIW, the article mentions the book "Designing Data-Intensive Applications" by Martin Kleppmann. I wanted to throw out my own endorsement for the book, it's been instrumental in helping me design my own fairly intensive data pipeline.
- jacobolus 9y ago6 days to go in the kickstarter map https://www.kickstarter.com/projects/1407076797/a-map-of-the-distributed-data-systems-landscape https://www.kickstarter.com/projects/1407076797/a-map-of-the...
- teacpde 9y agoNice, I thoroughly enjoy those maps at the beginning of each chapter.
- sunnykgupta 9y ago+1 for DDIA by Kleppmann!
- dswalter 9y agoIt's such a wonderful book. Reading it pushed me from thinking in terms of what I had worked with to building systems based on what was needed. I cannot recommend it highly enough, for pretty much anyone in {frontend, backend, data science, etc}.
- pcsanwald 9y agoI third this recommendation. I've worked on a ton of data intensive applications on all kinds of stacks over the years, and this book gives you lessons learned as well as a very valuable historical perspective on relational databases that is missing from a lot of the popular literature today.
- teej 9y agoDear HN reader - if you're not quite ready to buy the book, take a listen to this episode of Software Engineering Daily (https://softwareengineeringdaily.com/2017/05/02/data-intensive-applications-with-martin-kleppmann/ https://softwareengineeringdaily.com/2017/05/02/data-intensi...). It will give you a sense of what Martin Kleppmann is all about and how he thinks about problems. I ordered my copy of "Designing Data-Intensive Applications" after listening to this episode.
- mateuszf 9y agoIsn't this just event sourcing?
- cturner 9y agoThis is a flawed architecture. It will work at release, but it will be difficult to manoeuvre with, and they will grow to hate it. As your business changes, your data changes. Imagine if on day one, they had one author per article. On day 1000, they change this to be a list of authors. Kafka messages are immutable. Each of those green boxes on the right hand side of the first diagram will need to have special-case logic to unpack the kafka stream, with knowledge of its changes (up until 17 May 2017, treat the data like this, but between then and 19 May 2017 do x, and after that do y). Document pipelines is a rare instance of a context where XML is the best choice. They should have defined normalised file formats for each of their data structures. Something like the gateway on the left of the first diagram would write files in that format. (At some future time, they will need to modify the normalised formats. Files are good for that. You can change the gateway and your stored files in coordination.) Secondly, they should have a gateway coming out of the file store. For each downstream consumer, they should have a distinct API. These APIs might look the same on the first day of release. But they should be separate APIs so that you are free to refactor them independently. When you have a one-to-one API relationship, you can negotiate significant refactors in a single phone call. When you have more than one codebase consuming, you need to have endless meetings and project managers. I call this, "The Principle of Two." Some of the other comments here say that they should have used databases. So far, they have not made the case for it. And databases are easily abused in settings like this one. People connect multiple codebases to them, and use SQL as a chainsaw. Again, you can't negotiate changes. When you create a system, your data structures are the centre of that system. You need to do everything you can to keep your options open to refactor them at a later time, and to do so in a way that respects APIs that you are offering your partners. Kafka is a good tool. If used well, your deployment design will stop your system regularly (e.g. every day), nuke the channels, recreate them from scratch, and restart your system against these empty channels. You shouldn't use it as a long-term data store.
- elnygren 9y ago> Kafka messages are immutable. Each of those green boxes on the right hand side of the first diagram will need to have special-case logic to unpack the kafka stream, with knowledge of its changes (up until 17 May 2017, treat the data like this, but between then and 19 May 2017 do x, and after that do y). One solution would be: Kafka allows you to easily create new streams from the "monolog" stream that normalise the data to a certain schema. Consumers can then just consume these new derived streams. Another was mentioned in another reply (create a new stream that ultimately replaces the monolog). > Document pipelines is a rare instance of a context where XML is the best choice. XML does not really offer anything here that could not be achieved with tools that are nicer to work with? It's just a file format, basically. Why wouldn't Protobuf work? It also saves a huge amount of disk space vs. XML. > Secondly, they should have a gateway coming out of the file store. For each downstream consumer, they should have a distinct API. This is Kafka. They can have distinct streams for the consumers since you can always derive new kinds of streams. > You shouldn't use it as a long-term data store. Why not? Kafka has support for infinite retention and Kafka has very strong guarantees about always writing data to disk and not losing a single event.
- aug_aug 9y agoFigure 3: The Monolog, containing all assets every published by The New York Times.
- toomim 9y ago> Traditionally, databases have been used as the source of truth ... [but] can be difficult to manage in the long run. First, it’s often tricky to change the schema of a database. Adding and removing fields is not too hard, but more fundamental schema changes can be difficult to organize without downtime. This argument sounds self-contradicting. Kafka doesn't let you change its schema at all! At least postgres gives you the option. It seems that the author is excited about having a single source of truth that doesn't change, and didn't realize that he could do that with a database, if he just never used the schema-changing features. Am I missing something? It seems like the author could be totally happy with a bunch of derived postgres databases sitting in front of a "source of truth" database, where he never changes the source of truth database's schema. Why use kafka?
- elnygren 9y agoPostgres could in fact be used here by creating an append-only table similar to: id, data (JSON field) However, Postgres doesn't have good support for creating derived append-only logs (=streams) from that table. Kafka has Kafka Streams and KafkaSQL. And producer+consumer APIs that are a good fit for NYTs use case. > Am I missing something? Remember that in addition to schema changes, NYT also wants to avoid row changes. One reason was that all search indices & other systems need updating too during a row change in a DB and this will lead to inconsistencies in large scale (sometimes some of these updates fail). IMHO the thing you are missing is the log based architecture where all databases are derived/materialised from the SSOT: the log.
- qaq 9y agoHowever, Postgres doesn't have good support for creating derived append-only logs ???? A one line trigger will give you derived append only log that is transactionally consistent.
- elnygren 9y agoso how does that one-liner transform data from log A to log B in realtime? are talking about creating a new table or materialised views or stored procedures or...?
- look_lookatme 9y agoThis is very similar to a normalized model in a relational database, with many-to-many relationships between the assets. In the example we have two articles that reference other assets. For instance, the byline is published separately, and then referenced by the two articles. All assets are identified using URIs of the form nyt://article/577d0341-9a0a-46df-b454-ea0718026d30. We have a native asset browser that (using an OS-level scheme handler) lets us click on these URIs, see the asset in a JSON form, and follow references. The assets themselves are published to the Monolog as protobuf binaries. When consuming this data do you have to programatically do relationship fetching on the client side or is eager loading/joins available in some way in Kafka? Additionally there seems to be a focus on point-in-time specific views of this data, but are you able to construct views using arbitrary values/functions? Let's say each article is annotated with some geo data, can you construct regional versions of these materialized views of articles at the Kafka level? If not it seems like you are pushing a fair amount of existing sophisticated behavior at the RDBMS level up into custom built application servers.
- iooi 9y ago> Because the topic is single-partition, it needs to be stored on a single disk, due to the way Kafka stores partitions. This is not a problem for us in practice, since all our content is text produced by humans — our total corpus right now is less than 100GB, and disks are growing bigger faster than our journalists can write. Before this line, the author mentions they also store images. There's no way that all their text + images is <100GB right? Something is inconsistent here.
- newforice 9y agoWhen they publish hateful garbage like the following, I have no interest in their supposed technical prowess: https://www.nytimes.com/2017/09/04/us/texas-storm-federal-aid-abbott-cruz.html https://www.nytimes.com/2017/09/04/us/texas-storm-federal-ai...
- sctb 9y agoCould you please stop violating the guidelines with gratuitous off-topic flamebait? https://news.ycombinator.com/newsguidelines.html https://news.ycombinator.com/newsguidelines.html
- newforice 9y agoI stated an opinion. You don't like it? Tough.
- pizzaman09 9y agoWhat does this have to do with the article?
- deleted 9y ago[deleted]